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.

Zephon's operator model. The stages a sample passes through: the 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:

examples/guide/pipelines/first_pipeline.py
"""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:

examples/guide/pipelines/torch_dataloader.py
"""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.