A GPU cluster is usually specified compute-first: pick the accelerator, the count, the network, and then buy "enough storage". That ordering produces two expensive mistakes. One is an all-flash parallel file system sized for a data-read rate the job never needs. The other is a storage system that looks idle on average and still stalls a thousand GPUs for minutes at every checkpoint and every restart. Compute and storage co-design means sizing the storage path from the job's actual traffic, burst by burst, and shaping the job so that traffic is cheap to serve.
Other articles on this site go deep on single paths: GPUDirect Storage, the checkpoint save path, DataLoader bottlenecks and KV-cache disk offload. This one is the whole-cluster view: the four traffic classes that compete for the same storage, which tier should serve each, a worked sizing of a 1,024-GPU training job, and the interference failures that only appear when everything runs at once.
Four traffic classes, one storage system
Every storage request from a training or serving cluster falls into one of four classes, and they have very different shapes:
| Class | Shape | What it needs |
|---|---|---|
| Training data reads | Steady, continuous, every rank; random access across shards | Sustained throughput with low tail latency; a few seconds of prefetch hides jitter |
| Checkpoint writes | Periodic burst; every rank writes its shard at the same moment | Burst write bandwidth, many concurrent creates, durable commit |
| Restart reads | Rare, all ranks at once, after a failure while the cluster is idle | Aggregate read bandwidth; every second is the whole cluster's time |
| Model and KV loads | Serving cold starts, autoscaling waves, cache offload | Fast fan-out of the same large files to many nodes |
The important property is that the expensive classes are bursty and correlated. Each rank's checkpoint shard is modest, but all ranks write at the same instant, and a restart is by definition a moment when every GPU waits on storage. Averages hide this: a storage system that is 2% utilised over a day can still be the reason a job loses an hour of cluster time to restarts that week.
Worked example: budgeting a 1,024-GPU job
Take a dense 70-billion-parameter model on 1,024 GPUs (128 nodes of eight). Assume each GPU sustains 400 TFLOP/s of useful bf16 math, that training costs about 6 FLOPs per parameter per token, that token ids are stored as 4-byte integers, and that the parallel file system delivers 50 GB/s aggregate reads. Assume too that the full training state is replicated in each of 16 data-parallel replicas; with a distributed optimizer only the weights replicate, and the restore amplification below shrinks. Every number below follows from those assumptions; swap in your own.
def storage_budget(gpus, params, tflops_per_gpu, bytes_per_token=4,
ckpt_bytes_per_param=14, dp_degree=16, read_gbs=50.0):
"""Back-of-envelope traffic for dense LLM pre-training. All inputs are assumptions."""
tok_s_gpu = tflops_per_gpu * 1e12 / (6 * params) # ~6 FLOPs per param per token
text_read_mb_s = tok_s_gpu * gpus * bytes_per_token / 1e6
ckpt_gb = params * ckpt_bytes_per_param / 1e9 # bf16 weights + fp32 master + Adam m, v
return {
"tokens/s per GPU": round(tok_s_gpu),
"text read MB/s (cluster)": round(text_read_mb_s, 1),
"checkpoint GB": round(ckpt_gb),
"sync save in 60 s, GB/s": round(ckpt_gb / 60, 1),
"async drain over 30 min, GB/s": round(ckpt_gb / 1800, 2),
"restore, every DP replica reads, s": round(ckpt_gb * dp_degree / read_gbs),
"restore, read once and broadcast, s": round(ckpt_gb / read_gbs),
}
print(storage_budget(gpus=1024, params=70e9, tflops_per_gpu=400))| Quantity | Value | Consequence |
|---|---|---|
| Tokens per GPU per second | ~952 | Data rate per rank is tiny |
| Pre-tokenised text reads, whole cluster | ~4 MB/s | Any storage tier can feed this; it is not the design driver |
| Image training instead (1,000 images/s/GPU at 150 KB) | ~154 GB/s | Now reads are the driver: cache shards on local NVMe |
| Full training state at 14 bytes/parameter | ~980 GB | Roughly a terabyte per checkpoint |
| Synchronous save inside 60 s | 16.3 GB/s | Every GPU idles for the whole write |
| Asynchronous drain over a 30-minute interval | 0.54 GB/s | Over 25 times less bandwidth, if staging memory exists |
| Staging per node if the 16 replicas split the save | ~7.7 GB | Fits comfortably in host DRAM |
| Restore where all 16 data-parallel replicas read | 15.7 TB, ~314 s | Read amplification by the DP degree |
| Restore: one replica reads, fabric broadcasts | ~20 s | Same bytes from storage once |
Three design facts fall out. For text pre-training, data reads are a rounding error, so spending on read bandwidth for them is waste. Checkpoint cost is a choice: synchronous saves demand over 16 GB/s, while staging to host memory and draining in the background needs well under 1 GB/s, so the same storage system can look either inadequate or oversized depending on the software. And restart time is dominated by read amplification, not raw bandwidth: the 16× factor costs more than any upgrade would buy back. Byte accounting for the checkpoint itself is covered in the checkpointing deep dive.
Which tier holds what
Co-design assigns each class to the cheapest tier that meets its burst, and makes the software move data between tiers deliberately.
- Host DRAM is the checkpoint shock absorber. Copy GPU state into pinned host buffers (seconds over PCIe), let training resume, and drain to slower tiers in the background. Budget the memory up front: staging plus DataLoader prefetch plus page cache must not push the node into swap or the out-of-memory killer.
- Local NVMe is the fast, cheap, non-durable tier. Use it for a cached copy of the dataset shards each node reads, a tier-0 checkpoint that makes restart after a software crash nearly free, and a model-weight cache for serving. A single PCIe Gen4 x4 drive reads at most about 7 GB/s sequentially and nodes usually carry several; treat those as ceilings and measure random-read throughput with your real record size.
- Parallel file system is the shared, durable, high-bandwidth tier. It holds the checkpoints you would restart from after losing a node, and datasets too large to replicate. Its scarce resources are aggregate bandwidth during bursts and metadata operations, not capacity.
- Object storage is the source of truth and the retention tier: cheap, durable, high aggregate throughput with many parallel streams, but with per-request latency and request-rate limits. Amazon S3, for example, documents at least 3,500 write-type and 5,500 read-type requests per second per partitioned prefix, so spread keys across prefixes when thousands of ranks hit it at once.
Data should move down this ladder on a schedule the software controls: object store to local NVMe before the job starts, local NVMe into host memory just ahead of use, GPU to host memory at checkpoint time, then host memory to file system and object store at leisure. When the job instead pulls from the slowest tier on the critical path, the storage system sets the step time. Deterministic, replayable data ordering, described in the training data pipeline article, is what makes caching and exact resumption safe at the same time.
Shaping the job to fit the storage
The cheapest storage upgrade is usually a change to the job, not to the hardware.
- Pack small files. Millions of individual images or JSON files turn every epoch into a metadata storm. Pack them into shards of hundreds of megabytes (tar-style WebDataset shards, Parquet, or a tokenised binary with an index) so reads are large and sequential and the file count stays in the thousands.
- Read once, broadcast. On restore, have one data-parallel replica, or one rank per node, read each shard and send it over the GPU fabric, which is usually far faster than storage. The same applies to serving cold starts: pull weights once per node into local NVMe or host memory, then load every GPU from there.
- Stagger and rate-limit bursts. Spread checkpoint drains across ranks so they do not open thousands of files in the same millisecond, and cap drain bandwidth so background writes never starve foreground reads.
- Bypass the page cache for large streams. Buffered reads of checkpoint files copy every byte through the kernel page cache, a user buffer and a pinned buffer before it reaches the GPU, and evict the dataset cache on the way. Direct I/O, or GPUDirect Storage where the stack supports it, removes copies and protects the cache.
- Shard for the fabric. Write checkpoints in a layout that allows resharding on restore, so a job that restarts on a different GPU count does not have to read everything onto every rank.
Verify every one of these changes with measurements taken from the job, not from a storage benchmark. Instrument three numbers per step: time the training loop spends waiting for the next batch, time blocked inside the checkpoint call, and wall-clock time from failure detection to the first completed step after restart. Then reproduce the burst synthetically: a benchmark such as fio, run from every node at once with the real file sizes, block sizes and number of open files, tells you what the storage does under the checkpoint and restart patterns. A single-client benchmark measures the drive or the network link, which is rarely the limit; the concurrent version measures metadata servers, switch uplinks and the storage controllers, which usually are. Repeat the concurrent test while a second job is running, because a shared file system has neighbours.
Interference failures
- Checkpoint burst starves data reads. Symptom: step time spikes for a few minutes after each save even with async checkpointing. Cause: the background drain saturates the same NICs or file-system servers the data loader uses. Fix: rate-limit the drain, cache data on local NVMe, or put storage traffic on a separate network.
- Restart storm. After a failure, every rank stats, opens and reads the same checkpoint. Metadata servers saturate before bandwidth does. Fix: one reader per shard, broadcast, and fewer, larger files.
- Silent partial checkpoint. A drain fails halfway and the job records the checkpoint as done. Fix: write to a temporary path, verify sizes or checksums, then commit with an atomic rename or manifest, and never delete the previous good one first.
- Host memory exhaustion. Async staging, pinned prefetch buffers and page cache grow together until the kernel kills a rank. Fix: a written per-node memory budget and pre-allocated, reused staging buffers.
- Local NVMe wear and loss. Tier-0 checkpoints on node-local drives vanish with the node and wear the flash with constant rewrites. Keep them as an accelerator, never the only copy.
- Serving autoscale stampede. Twenty new replicas each pull 140 GB of bf16 weights for a 70B model from the object store at once. Fix: node-level weight caches, pre-warmed nodes, and peer-to-peer distribution.
Trade-offs
| Decision | Option A | Option B |
|---|---|---|
| Checkpoint path | Synchronous to shared FS: simple, needs huge burst bandwidth | Async via host DRAM: tiny bandwidth, needs memory and careful commit logic |
| Dataset location | Stream from shared FS: no staging step, shared bottleneck | Cache on local NVMe: fast and isolated, needs a pre-stage and capacity per node |
| Storage network | Converged with the compute fabric: fewer cables, interference risk | Separate storage network: isolation, more cost |
| Checkpoint frequency | Frequent: less lost work per failure | Rare: less I/O, more lost work; price both using the failure rate |
The fabric itself is part of the storage design, since broadcasts and drains share it with collectives; the InfiniBand article covers that layer.
What to do next
- Fill in the budget function with your model size, GPU count, measured throughput and data format, and write down the four traffic numbers.
- Measure your storage under a burst that matches the checkpoint (all ranks, real file sizes), not a single-client benchmark.
- Switch to asynchronous checkpointing with an atomic commit, and budget host memory for it.
- Make restore read each shard once and broadcast it; time a full restart end to end.
- Pack small files into large shards and cache them on local NVMe if reads exceed a few GB/s.
- Rate-limit background drains and watch data-loader wait time around each checkpoint.