Datology AI
Open Source Release18 min readOctober 2026

Zephon: Fast, Flexible, Elastically Deterministic Data Loading

Zephon: Fast, Flexible, Elastically Deterministic Data Loading hero image

If you want to know whether a change to your training data made a model better, everything else about the training run has to hold still. In practice, the data loader often doesn't: resume a run on a different number of GPUs, and most loaders will quietly change what each GPU sees, which in our experiments has been enough to muddy comparisons between data curation methods.

At Datology, we run thousands of data experiments, and more and more of them are designed and run by agents, which (much like us humans) won't notice when a result depends on the GPU count.

Today, we're open-sourcing Zephon, our data loader for text and multimodal training. Zephon preserves the same sequence of global training batches across changes in GPU count, operator parallelism, and execution backend, even for pipelines that tokenize, pack, and mix data during training.

Data prep on the fly

Figure 1: Training a foundation model consists of offline and online data preprocessing steps. Today’s typical pipelines shift many more things to the offline stage than needed, often storing the fully processed input tensors on disk.

Whether to materialize a view on your data or keep it virtual is a choice databases have long studied. A materialize-versus-recompute trade-off exists in many systems, from Spark's lineage-based recomputation to activation checkpointing in LLM training. Materializing data makes sense when repeated use justifies the upfront processing and storage costs, or when online processing would not be fast enough to keep a consumer such as a GPU from waiting.

However, in the context of training foundation models, the cost of offline materialization is often underappreciated: each change to the data pipeline can require another full pass, processing time, and another stored copy. It also barely amortizes. Training runs are often single-epoch, so for many ablations the stored corpus is read exactly once — and a corpus read once never earns back the cost of writing it.

For some data modalities, such as video, materializing tensors of the fully decoded training data is fundamentally infeasible due to the massive inflation in size. Labs that work with such modalities resort to building proprietary loading infrastructure.

To avoid the storage and time costs of materializing training datasets offline for every new training job, we want to move more of the data preparation into the training job itself. The catch is that offline materialization is not just convenient; it hands us two things almost for free. It is fast at training time, because a preprocessed corpus just needs to be loaded from disk. And it makes determinism straightforward, because the final dataset on disk has a clear index of which sample is used when in training - more on that below. Zephon is our solution for offering fast and deterministic data loading while processing data online.

Your data loader is a hidden variable

If you resume a training run on 8 GPUs instead of the 16 GPUs it started on, you would expect to be training on the same data in the same order. That is not the case. The data loader quietly changes what each GPU sees. Mosaic Streaming coined the property you might have been expecting elastic determinism: yielding the same global batches no matter how many GPUs you train on. Most data loaders do not even offer it.

Isolating variables is the basis of experimental research, as early systems work stressed. If two runs differ in one parameter and everything else behaves deterministically, the difference in outcome belongs to that parameter. This is not the case when the loader silently changes the data order, making it a hidden variable in your experiments.

Figure 2: When training resumes after a fault or cluster resize, a standard loader’s loss curve and batch sequence depend on the number of GPUs. Zephon preserves both across GPU counts. Drag the divider to compare; each loader is shown against its own uninterrupted reference.

Any variation in the input data sequence can confound the results of ablations. To illustrate this, we train a 1B-parameter model on 16 GPUs, take a checkpoint at step 150, and then resume from that checkpoint on 8, 32, and 64 GPUs alongside the original run, once with a standard data loader and once with Zephon (think of it as the run hitting a hardware fault and coming back on whatever set of GPUs happened to be free).
As seen in Figure 2, with the standard data loader, changing the GPU count means that the data loader hands the model a different sequence of samples after the resume, and so each one ends up with a distinct loss curve. With Zephon, the curves lie on top of each other, and the SHA-256 hash of every global batch matches the baseline run.

To get a sense for how much this matters, we train the same 1B-parameter model from scratch for 20B tokens on 8, 16, 32, and 64 GPUs using a standard data loader. We show the scores per GPU count in Figure 3 and Figure 4. The model’s score on the FineWeb evaluation suite varies by up to 0.64 points depending on how many GPUs it was trained on, and its DCLM Core v1 score by up to 0.82 points. That is roughly four times the thresholds that the FineWeb and DCLM authors used when evaluating data curation methods, which means that an ablation comparing two mixtures on two different allocations could easily be measuring the GPU count instead of the data. With Zephon, the spread across the different GPU counts was only 0.011 and 0.05 points, respectively.

Additional contextLearn more about the eval thresholds

FineWeb is Hugging Face’s 15-trillion-token English web dataset. They used controlled training experiments to select filtering and deduplication methods. Its evaluation suite averages scores across eight benchmark groups selected for their signal at small training scales: CommonsenseQA, HellaSwag, OpenBookQA, PIQA, Social IQA, WinoGrande, ARC (Easy and Challenge), and MMLU. The relevant threshold we identify for identifying differences between ablations is 0.15 points. This is the smallest estimated gain among three custom filters adopted after 28 B-token training experiments: the short-line filter scored 42.65%, compared with an approximately 42.50% baseline. We reconstruct this difference from the endpoint in the FineWeb Report, Appendix E.4, Table 2, and the baseline in Figure 7; this illustrates an improvement that informed an actual filtering decision.

DataComp-LM (DCLM) is a benchmark for comparing language-model training data under controlled training recipes; its associated CORE suite is its aggregation of evaluations that they use to measure whether filtering, deduplication, and data mixing improve downstream performance. CORE averages chance-normalized scores over 22 task configurations: ARC-Easy, ARC-Challenge, BoolQ, CommonsenseQA, COPA, CoQA, HellaSwag (0- and 10-shot), Jeopardy, LAMBADA, OpenBookQA, PIQA, SQuAD, Winograd, WinoGrande, AGIEval LSAT-AR, and six BIG-Bench tasks (QA Wikidata, Dyck Languages, Operators, Repeat Copy Logic, CS Algorithms, and Language Identification). Our comparison value of 0.2 points comes from the authors’ deduplication study comparing Bloom filtering and MinHash plus Suffix Array, where the authors say results within this range are comparable.

These comparison values come from published experiments; they are not statistical significance thresholds. Our results show that changing GPU topology can introduce score differences large enough to complicate comparisons between data curation methods, potentially masking improvements or making one method appear better than another.

Training a model on a varying number of GPUs is not as exotic as it might sound at first. GPU resources are scarce. It is standard practice to prototype runs intended for large GPU counts on smaller allocations, using gradient accumulation to simulate the target global batch size. Moreover, cluster priorities may shift throughout the life of a training job. A more urgent job may preempt another job, forcing it to resume on fewer nodes. Conversely, when additional nodes become available, a job should be able to scale up to finish sooner. In all these cases, the data loader should produce identical global training batches regardless of the GPU topology.

The data loading community has been looking into this problem, but the existing solutions assume indexable datasets or fixed operator parallelism.

Where existing data loaders fall short

We love Mosaic Streaming. We relied on it internally before building Zephon. However, modern foundation model input data pipelines consist of stateful operations with n-to-m characteristics. For example, in pre-training, we might split documents into multiple subsequences after tokenizing them, pack these subsequences into stateful bins, and reorder samples to guarantee the correct mixture. Streaming's guarantees don't extend to such pipelines.

Figure 5: Online processing of samples can, due to operations such as packing, distort the ordering of samples that we see. Such stateful n-to-m pipelines do not have a sample index that represents the sample to train on at a certain step. In this animation, you can see that an input sample can correspond to multiple output samples, and that inputs don't have to show up in the same order in the output.

Whereas in a 1-to-1 pipeline, progress is a sample index and resuming means looking up the last accessed sample, in a stateful n-to-m pipeline, logical progress is spread across in-flight operator state, partially formed outputs, and reordered records. These pipelines are not indexable. The question “Which sample will we train on at step 1234?” cannot be answered in a non-indexable pipeline.

Mosaic Streaming and Grain provide elastic determinism for indexable pipelines. Megatron Energon supports online packing and reproducible scaling, but requires a fixed total number of workers. Its recovery relies on snapshots of operator state and replayable transforms to reconstruct buffered samples. We wanted to preserve identical global batches for stateful n-to-m pipelines while changing GPU count, operator parallelism, and execution backend. So we built Zephon.

Zephon

We built Zephon to be modular. No matter if you want to use Parquet, Vortex, or LitData, or if you want to tokenize offline but pack online, or if you just want to load pretokenized, prepacked data - you can use the system in all these scenarios and rely on Zephon’s built-in operators. Of course, Zephon also supports custom user-defined operations.

Zephon is available as a Python package. Zephon pipelines are defined via a builder-pattern API (Code block 1):

ws = StaticMixtureWorkSource(
    [Dataset.from_path("demo", "s3://example-dataset")],
    mixture=MixtureSpec({"demo": 1.0}), chunk_size=1024)
pipeline = (Pipeline(ws)
    .tokenize(tokenizer_id="gpt2", parallelism=4)
    .pack_sequences(target_length=2048)
    .ensure_mixture()
    .batch(microbatch_size=8))
for batch in pipeline:
    train_step(batch.to_training())

Code block 1: An example of Zephon's builder pattern API.

Lines 1 to 3 define the training curriculum via a WorkSource. Lines 4 to 8 attach a chain of operators that turn raw data into training batches.

The WorkSource decides what to train on: which datasets to use, in what proportions, and in what order. It produces sample pointers, not the data itself. A pointer identifies a sample within a dataset, so the curriculum can be defined without reading all of the sample payloads.

Zephon’s fetch operator (implicitly the first in every pipeline) resolves those pointers into actual data payloads. The records then flow through whatever operators you defined, finally reaching the training loop as batches. The WorkSource and the processing pipeline are separate, so changing the data curriculum does not require rewriting the preprocessing code, and optimizing an operator does not require changing the curriculum.

In addition to elastic determinism for stateful n-to-m pipelines, Zephon comes with a wide variety of additional features:

As curriculum is separated from processing:

  • Token-aware mixing: Data mixing ratios are typically defined on the token level. Now what if your data lives as strings in your input files? Zephon’s got you - the system measures the composition of your datasets and yields the data such that post-tokenization your model sees the actual intended mix.
  • Dynamic curricula: Zephon decouples the curriculum from downstream operator processing. This allows researchers to work on dynamic data mixes while engineers optimize the downstream processing pipeline. Keep an eye open: we will release an easily configurable component for curriculum learning soon!
  • No format lock-in: The curriculum is defined over sample pointers, not payloads, so the file format only matters when fetching. Zephon ships with optimized implementations for Parquet, LitData, MDS, Vortex, and JSONL, and another format is just a future PR away.

Due to our flexible operator design we can offer:

  • Packing on the fly: Zephon ships with a variety of packing operators such as wrap or bin packing that work across modalities.
  • Online tokenization: For pre-, mid-, and post-training, Zephon offers a flexible tokenization operator that can be configured to your needs.
  • Custom operators: Write your own functions and full operators whenever you need your own custom transformations, batching, or whatever else is required for your data pipeline.
  • Cloud-native: Stream your data from S3, GCS, or Azure with optional prefetching that warms a local cache ahead of consumption. Local and HPC distributed filesystems work out of the box as well.

Because a pipeline is just an iterable:

  • Framework-agnostic: Your data pipeline shouldn't be coupled to your training framework. Zephon pipelines are Python iterables that yield numpy- or torch-ready batches. Use them with PyTorch or whatever you train with.

Before we get into how Zephon works, we want to show you that we don’t compromise on speed. After that, we will walk through what it took to keep our guarantees.

Keeping GPUs fed

A data loader must not bottleneck training. This means that it needs to deliver batches faster than GPUs can consume them in a given training setup.11.Additional throughput beyond this has no impact on downstream training, but may matter in different setups, e.g., a smaller model may need higher throughput as the GPU can do the fwd/bwd passes faster. The required rate of data loader throughput cannot be generally stated and is a function of the hardware, training framework, and model architecture. As discussed, when moving work out of the offline stage into the online training job, the loader needs to keep up with the GPUs.

To start off, we measure Zephon’s ability to just move data: fetching and batching. In this pipeline using VLM training data (Figure 6), Zephon reaches 16k samples/s. The strongest baseline, Grain, reaches 5.8k. LitData and Streaming never scale past roughly 2k samples per second, at any worker count.

Additional contextInside the throughput benchmark

All benchmarks run on a dedicated AWS r8id.48xlarge instance: Intel Xeon 6975p /w 96 physical cores, 1.5 TiB RAM, 3x3.5 TB NVMe RAID-0, on Ubuntu 24.04.

Vision-language samples are multimodal documents from the MAmmoTH-VL dataset.

For fetch-and-batch, the data loading pipeline loads the raw sample payloads and creates batches, involving minimal processing.

For the full VLM pipeline, we rely on a non-monotone, standard VLM preprocessing pipeline. This pipeline reads raw records, decodes text and images, applies vision transforms, applies the chat template and tokenizes the sequence, shuffles, packs, and constructs batches.

Regarding parallelism, the baselines note the worker count, while for Zephon it depends on the benchmark. For fetch-and-batch we vary the degree of parallelism of the fetch operator, while for the full pipeline it is the sum of per-operator parallelism across all pipeline operators.

A test more representative of the real world is evaluating a full VLM pipeline (Figure 7), which includes decoding, transforming, tokenizing, packing, and batching. Here we report usable tokens per second per GPU, counting only non-padding tokens, because a loader that packs badly would otherwise look fast while feeding the model only padding tokens. The horizontal line marks the consumption rate of the model in our training setup; any system scoring above it keeps the GPUs busy.

Zephon and Grain are effectively tied at 32-way parallelism, and both clear the minimum throughput threshold. Both deliver more than 23 times the throughput of the Streaming-based solution we had quickly hacked together internally before we started working on Zephon. That loader plateaus as parallelism rises and never reaches the 16.9k usable tokens/s/GPU this model needs to train without stalling.22.It has to shoehorn image transformations, packing, etc. into streamings 1-to-1 model - effectively all operations are serialized, which is bad for throughput.

Note that Zephon reaches this throughput while guaranteeing elastic determinism for this stateful n-to-m pipeline, which none of the baselines provides. The guarantees are cheap enough that you do not have to choose between them and keeping your accelerators busy. For more results, including time to first batch and the simpler text pipelines, please check out our paper.

Keeping our guarantees

Having shown that Zephon is fast enough, we can turn to how we achieve both high throughput and the guarantees. Over the life of a training run with Zephon, there are a handful of knobs you’ll want to turn: how many workers each operator gets, whether those workers are threads or processes, how many GPUs you’re training on, and how you’ll want to resume the job after a failure or preemption. We designed Zephon so that none of these decisions changes the data your model sees. Before going through them one at a time, we need to say a little bit about how a Zephon pipeline is put together.

The most common way to parallelize a data loader, and the one that PyTorch’s DataLoader uses, is to give each worker its own copy of the entire pipeline. In this model, every worker reads data, transforms it, and produces finished batches. This has the virtue of simplicity, but it also means you get exactly one performance knob, and turning it adds parallelism to every step at once. If a particular transformation step (e.g., tokenization) is your bottleneck and you want to give it more CPU, you may inadvertently end up with a lot more processes trying to read data from your distributed filesystem at the same time.

Figure 8: Zephon’s operator model.

Zephon instead represents the pipeline as a graph of operators (Figure 8), and each operator can have its own degree of parallelism. This matters for example for multimodal pipelines, where decompressing and transforming images vs. tokenizing text have very different CPU and memory profiles. In Zephon, operators pass records to each other via bounded queues, so if image transformation falls behind, the queue in front of it fills up and the upstream operators slow down instead of building an ever-growing backlog in memory.

Once a pipeline is spread across many operators running in parallel, its state gets spread across them too, and keeping the output deterministic requires some care. The following subsections detail how Zephon handles each of the different knobs.

Determinism that survives parallelism

Every Zephon operator is split into two parts: an accumulator, which sees every record in order and makes all of the stateful decisions, and a stateless transformation that does the expensive work across however many workers you give it. For example, the packing operator’s accumulator keeps track of the partially filled bins and decides which sequences belong together, and the workers do the actual work of assembling the packed sequences. Workers can finish in any order, so Zephon numbers each group of records before handing it out, and a reorder buffer puts the results back in sequence before passing them downstream. With some additional care around RNG state, this makes parallelism a pure performance knob, so that you can double the number of workers and the output won’t change.

Figure 9: Zephon’s runtime architecture in action. We show the execution of the fetch operator. You can see how the accumulator receives inputs (pointers to samples), generates microbatches and sends them to the worker processes. The workers fetch the sample payloads and emit to the result queue. Because workers do not necessarily finish in the same order as they receive input, the reorder buffer restores the deterministic input order.

All of the built-in operators handle this automatically, and the map_transform operator allows you to do this for custom transformation and filtering functions. Fully custom operators have a few rules that they have to follow for state and randomness, which our documentation walks through.

Determinism that works across execution substrates

The operators run on a runner, which is the execution substrate that actually does the parallelizing. Zephon’s process runner lets CPU-bound Python code run in parallel on OS processes despite the Global Interpreter Lock (GIL), which is the standard approach in data loading. The thread runner keeps the workers inside of a single process and uses threads instead, so no IPC and Python serialization is needed. Threads are a good fit for I/O and for native code that releases the GIL, and on free-threaded Python versions they can run Python code in parallel too. The ordering machinery from Determinism that survives parallelism lives above the runner abstraction, so you can switch from one to another without rewriting your pipeline or changing its output. We think that the future of Python is GIL-free. We built Zephon to be ready for it.

This flexibility isn’t free: in the operator model with the process runner, intermediate results cross an IPC boundary between each pair of operators. Zephon makes this affordable by coalescing the large tensor payloads in each group of records into one shared memory buffer per dtype, so the OS pipe used by the mp.Queues only carries a small amount of metadata. The shared memory coalescing section of the paper describes this in detail.

Determinism that survives a change in GPU count

Everything we’ve discussed so far is designed to keep the output identical when running on a fixed number of GPUs. To keep it fixed when the GPU count itself changes, Zephon divides the data into a fixed number of logical streams called lanes instead of dividing it up among the GPUs directly. The WorkSource produces a separate ordered stream of sample pointers for each lane, and every stateful operator keeps separate state for each lane, so a lane produces exactly the same records no matter which GPU it is assigned to, and changing the number of GPUs just redistributes the lanes.

Figure 10: Zephon’s lane multiplexing guarantees identical global batches across different GPU counts. Lanes are deterministic as Zephon’s operators are deterministic, the global batch is deterministic due to the way we yield data from each lane.

To turn identical lanes into identical global batches, at the end of a pipeline, Zephon’s lane multiplexer hands out training microbatches round-robin across the lanes each GPU owns, holding back faster lanes until the slower ones catch up so that every lane contributes the same amount of data to every optimizer step. In Figure 10, a global batch of 16 samples comes from four lanes: on two GPUs, each GPU accumulates gradients over one microbatch from each of its two lanes, and on four GPUs, each processes one microbatch from its single lane. Either way, the global batch contains the next four samples from every lane. This works as long as the lanes divide evenly across the data-parallel replicas and each replica’s gradient-accumulation steps cover whole rounds through its lanes.

The number of lanes has to be chosen up front. A common multiple of the data-parallel sizes you expect to run is a good choice: 32 lanes cover 8, 16, and 32 GPUs. Choosing far more lanes than needed isn’t free, since every additional lane per GPU increases the time it takes to restore from a checkpoint.

Your job will get preempted

As we mentioned in Where existing data loaders fall short, a stateful n-to-m pipeline has no sample index, so progress isn’t a number you can just write down. Instead, it’s spread across partially filled bins, shuffle buffers, and records that are still being reordered. So what goes into a Zephon checkpoint?

Zephon does not save operator state in its checkpoints. Instead, it records which chunks of sample pointers were still in flight when the checkpoint was taken, and the bounded queues keep that list short. Replaying only those chunks isn’t sufficient on its own, because an accumulator’s state can depend on a chunk that finished processing a long time ago. So for pipelines with stateful operators like packing, Zephon periodically sends a flush sentinel through the pipeline, every eight chunks per lane by default. The sentinel tells each accumulator to emit whatever it’s holding and start over from a clean slate. When a run resumes, Zephon replays the input from the most recent flush, which rebuilds the operators’ state exactly, and throws away anything that the training loop has already consumed.

Flushing adds a small amount of padding, about 0.02 percentage points of tokens at the default settings, and in exchange the checkpoints stay small and bounded. Together with lanes, this is what lets you resume a run on a different number of GPUs and get the same global batches. Once again, the paper has the full details here.

Why are you open-sourcing Zephon, and what’s next?

We originally built Zephon for our own data curation experiments at DatologyAI. Even with the help of the strong coding agents we have access to today, getting Zephon’s primitives right was tricky and took longer than we expected; elastic determinism for stateful pipelines has a number of subtle failure modes. The operator-accumulator split, lanes and flush sentinels are each what we arrived at after the simpler options fell over: giving every worker a copy of the whole pipeline, resuming from a sample index, and checkpointing operator state directly. We learned a lot in the process. We want to help create a world in which any team can build its own foundation model, so we decided to share what we built and learned.

This is a first version, and there are a number of ideas we’re excited to work on next. We are currently working on support for dynamic curricula, where the WorkSource changes the mixture weights over the course of a run while still keeping the determinism guarantees described above. We have also started prototyping a Ray-based execution engine (you can see traces of this in the repository already), which is relevant for compute-intensive audio and video preprocessing. We would like to build a more declarative interface to your data, so you don’t have to rely on the filesystem for defining mixes, and in tooling that makes it easier to tune pipelines so as to maximize their performance.

Most of all, we are looking forward to your feedback and hearing what the community might be interested in. If there is something that would be particularly exciting for you, please let us know!

Get Started

You can visit zephon.io, Zephon’s documentation, and Zephon’s repository on GitHub today! Our research paper explains the system in much more detail.

To help you get going, we provide example integrations of Zephon into two commonly used training frameworks: Torchtitan by Meta and the Pytorch Foundation, and Megatron-LM by NVIDIA. These examples are great starting points to integrate Zephon into your own codebase, or you can just clone one of them and get going immediately! The Documentation we mentioned above also walks through the integration into the training system.

Citation

If you use or refer to Zephon, please cite it as:

Maximilian Böther, Josh Wills, Ties Robroek, Sonnet Xu, Paul Burstein, Daniel Zayas, Cody Blakeney, Siddharth Joshi, Haoli Yin, Rishabh Adiga, Haakon Mongstad, Luke Merrick, Pratyush Maini, Ari Morcos, Matthew Leavitt, Ana Klimovic, and Bogdan Gaza. Zephon: Elastic Determinism for Online, Stateful Foundation Model Data Loading Pipelines. arXiv preprint, 2026.

Or use the BibTeX citation:

@article{Boether2026Zephon,
  author = {B{\"o}ther, Maximilian and
            Wills, Josh and
            Robroek, Ties and
            Xu, Sonnet and
            Burstein, Paul and
            Zayas, Daniel and
            Blakeney, Cody and
            Joshi, Siddharth and
            Yin, Haoli and
            Adiga, Rishabh and
            Mongstad, Haakon and
            Merrick, Luke and
            Maini, Pratyush and
            Morcos, Ari and
            Leavitt, Matthew and
            Klimovic, Ana and
            Gaza, Bogdan},
  title = {{Zephon}: Elastic Determinism for Online, Stateful Foundation Model Data Loading Pipelines},
  journal = {arXiv preprint},
  year = {2026}
}

Share this post

Share on TwitterShare on FacebookShare on LinkedIn

Ready for better data?

Let’s make models better through better data, automatically.

Book a Call