Pipelines: From a Curriculum to Training Batches¶
A Pipeline turns the curriculum defined by a WorkSource into the batches of tensors
that are consumed by your training loop. It fetches the underlying records defined in the
Datasets that comprise the WorkSource, applies any transformations that are defined on
the Pipeline to those records, and then optionally groups the resulting records into a
batch for training. For some workloads, your Pipeline definition may be very simple — it
may only need to fetch and batch records. For other workloads, you can define rich
transformations on your Pipeline in order to decode, transform, tokenize, shuffle, and
pack the samples prior to sending them to the training loop.
WorkSource turns a training
curriculum into mixture-aware per-lane chunks, the operator graph fetches and
transforms them, and the training loop — outside Zephon — consumes the result.
Building off of the web/code StaticMixtureWorkSource introduced previously, your first
pipeline might look something like this:
"""Fetch and batch records from a StaticMixtureWorkSource."""
from zephon import Pipeline
from zephon.io import Dataset
from zephon.work import MixtureSpec, StaticMixtureWorkSource
code = Dataset.from_path("code", "s3://my-bucket/corpora/code")
web = Dataset.from_path("web", "s3://my-bucket/corpora/web")
work_source = StaticMixtureWorkSource(
datasets=[code, web],
mixture=MixtureSpec({"code": 0.25, "web": 0.75}),
seed=42,
)
# Fetch is implicit, so this pipeline only groups the records it reads -- the
# right shape when the samples were prepared offline.
pipeline = Pipeline(work_source).batch(microbatch_size=32)
for sample_batch in pipeline:
print(len(sample_batch.records), "records")
This kind of simple pipeline definition is appropriate when the data you are training on
has mostly been prepared offline. The batch operator groups records into SampleBatch
objects, which Building Training Batches shows how to
convert into tensors for your model.
Most training frameworks are happy to consume data from any kind of Python iterable, like
a Zephon Pipeline, directly. However, if you are using a PyTorch-based training
framework that really requires a torch IterableDataset or a DataLoader, Zephon
provides an adapter method that allows you to treat a Pipeline object as if it were an
IterableDataset that can be passed in to a DataLoader like this:
"""Hand a Pipeline to a torch DataLoader that insists on a Dataset."""
from torch.utils.data import DataLoader
from zephon import Pipeline
pipeline = Pipeline(work_source).batch(microbatch_size=32)
# The Pipeline has already batched and already runs its own workers, so the
# DataLoader must do neither.
loader = DataLoader(pipeline.to_torch_dataset(), batch_size=None, num_workers=0)
In general, we do not recommend adding this extra level of conversion if you don’t absolutely need it; it’s generally simpler to act on the abstractions that Zephon provides for processing records and checkpointing state directly instead of trying to divide these responsibilities between the Zephon and PyTorch libraries.
The rest of this section takes you through the life of a Pipeline, from definition to
production. We start with an overview of how Zephon executes Pipelines, then cover the
built-in operators for turning raw samples into training data (preparing, shuffling,
mixing, packing, and batching records), and then discuss what you need to know to run a
Pipeline in a real training job, from performance tuning to distributed training and
checkpointing. Finally, we’ll close out this section with how you can write your own
operators for pipelines that need more than the built-in ones that we provide.