Storage & Data Pipelines
Keeping GPUs fed and their work safe: what training reads and writes, the data loader that starves a node, storage from NVMe to parallel filesystems and S3, checkpoint maths, GPUDirect Storage and model cold starts.
An interactive AI Infrastructure lesson: 20 steps, about 30 minutes, on a live simulation in your browser.
The team's 70B pre-training run, llama70, takes three of the four servers: 24 GPUs, tensor parallel 8 inside each server, ZeRO-3 across them. Its dataset lives on fs, the cluster's parallel filesystem: 40 GB/s read, 20 GB/s write, 1 PB.
Storage has to supply what the GPUs consume. These consume 22,964 tokens a second. The dataset is pre-tokenised, 4 bytes per token, so the run reads about 92 KB/s. A USB stick could keep up.
What you will learn
What a training job moves
- What an LLM run reads: The read rate a training job needs is its consumption rate times the bytes per sample: tokens/s × 4 bytes for pre-tokenised text, images/s × ~110 kB for JPEGs.
- The heaviest thing it writes: A resumable checkpoint is about 14 bytes per parameter (bf16 weights + fp32 master weights + two Adam moments): 7 × the size of the model you would serve.
- Datasets as shards: Pack training data into large shards: storage serves big sequential reads well and millions of small files badly, and shuffling shards is almost free.
Starving the GPUs
- Break it: a vision job on NFS: Data wait is the time the GPU spends with nothing to compute because the next batch is not ready. Prefetching hides loading only when loading is faster than the step.
- What starvation looks like: Low SM activity with a healthy GPU and network means the GPUs are waiting for input. Then ask whether the CPUs (workers) or the storage is the wall.
- Three times the workers: An input pipeline is only as fast as its slowest stage: storage bandwidth and CPU preparation are limits in series, and fixing one exposes the other.
- Fast storage and enough workers: Size the input pipeline from the GPUs' appetite: images/s the GPUs can train on, divided by images/s per worker, and times bytes per image for the storage.
- Drill: bandwidth to feed a node: Bandwidth to feed a node = GPUs × samples/s per GPU × bytes per sample. Do it before choosing storage, and multiply by the nodes that read at once.
Where the data lives
- Four places data can live: Local NVMe is fast and private, NFS is simple and slow, a parallel filesystem is fast and shared, object storage is huge and cheap but only fast for large sequential reads.
- Parallel filesystems and caches: Keep the master copy where it is cheapest and cache the working set where it is fastest. The GPUs only ever read the cache.
Checkpoints
- Checkpoint every 20 steps: Checkpoint overhead = checkpoint time ÷ time between checkpoints. At 49 s per 913 s it is 5.4 %, spent on insurance.
- Checkpoints to the NFS server: Checkpoint time = checkpoint size ÷ write bandwidth, and the GPUs idle for all of it unless the checkpoint is asynchronous.
- Break it: a crash between checkpoints: A failure costs the work since the last checkpoint plus the restart: relaunch time and checkpoint size ÷ read bandwidth.
- How often to checkpoint: Checkpoint interval ≈ √(2 × checkpoint cost × MTBF). As jobs grow, MTBF falls, so checkpoints must get cheaper to stay affordable.
- Asynchronous and sharded checkpoints: Asynchronous checkpointing turns a storage-bound pause into a seconds-long copy to host memory; sharding spreads the write over every node.
- Drill: checkpoint time: Checkpoint seconds = parameters × 14 bytes ÷ write GB/s: 70 B × 14 = 980 GB, ÷ 20 = 49 s.
GPUDirect Storage and cold starts
- GPUDirect Storage: GPUDirect Storage lets storage DMA straight into GPU memory, skipping the bounce buffer in host RAM: fewer copies, less CPU, more bandwidth per node.
- How long to load a model: Cold start = startup time + weights ÷ read bandwidth, and replicas starting together share the bandwidth. It sets how fast an inference service can scale up.
Recap & playground
- Cheat sheet
- Playground