Est.

High-Throughput Checkpoint Design for LLM Training at Scale

Failures at massive scale demand checkpoints so frequent that I/O overhead dominates training time.

Features Editor · · 11 min read
Cover illustration for “High-Throughput Checkpoint Design for LLM Training at Scale”
Systems Design · October 1, 2026 · 11 min read · 2,568 words

At extreme GPU counts, failures happen once per hour or more, and restart costs can come to dominate wall-clock training time, pushing systems into what one team calls a "restart-dominant regime". That's the anchor figure for this entire discussion: once a job crosses a certain scale, checkpoint overhead eats a material share of total training time, with documented cases showing it consuming a significant portion of wall-clock. The mechanism behind that is simple arithmetic. Mean time between failures shrinks inversely as GPU counts grow, so a job has to checkpoint more often just to bound its recomputation risk, while the volume of state it has to write keeps climbing alongside model and optimizer size. Those two curves move in the same direction, and that is what turns checkpointing into a first-order design problem rather than a footnote in a training run's postmortem.

The numbers from independent groups converge on the same story. One analysis frames checkpoint-restore itself as a big-data problem carrying all three of the classic "Vs": volume, in the form of terabyte-scale shards; variety, in the form of heterogeneous tensors of wildly different shapes scattered across thousands of ranks; and velocity, since some regimes checkpoint on every iteration. Another team measured checkpointing consuming roughly a tenth of total training time under an intentionally I/O-stressful per-iteration regime, with the share climbing well past that in more extreme configurations. The traditional response, checkpoint every N steps and restart from the last good state, pays for this twice over: once in I/O overhead paid on every single interval regardless of whether a failure ever occurs, and again in recomputation cost paid after each failure actually lands.

This is a storage architecture problem: how state moves from GPU memory down through host memory, local disk, and shared file systems sets whether a checkpoint costs seconds or minutes, and at 100,000-GPU scale that gap compounds into days of lost training time over the life of a run.

Checkpoint contents and write cost

A checkpoint is not one file. Model weights, optimizer states (first and second Adam moments), RNG state, and scheduler state are all sharded across ranks via 3D parallelism, tensor, pipeline, and data parallelism, so no single rank holds a coherent checkpoint, and the full state only exists when all ranks have flushed consistently. The full checkpoint only exists as a logical construct, assembled once every rank has flushed its shard consistently. That constraint alone explains why checkpoint I/O resists easy optimization: the thing being written doesn't live anywhere until the write itself completes.

Scale makes the shape of that data worse, not just its volume. Research characterizing production-scale checkpoint I/O describes a large number of heterogeneous tensors of different shapes and sizes spread over thousands of GPUs, grouped into shards and persisted across a correspondingly large number of files. Optimizer state usually dominates the byte count. In ZeRO-sharded training, the second moment tensor alone can outweigh the model weights themselves in raw size. A checkpoint's I/O profile is driven less by the model architecture than by which optimizer state the training run happens to be carrying.

Then there's the barrier problem. Synchronous checkpointing forces every rank to wait on the slowest write in the entire cluster before any rank can proceed, so one congested storage path or one slow node stalls thousands of GPUs simultaneously, on every single checkpoint interval for the duration of the run. Measurements of uncoalesced, small-buffer write patterns show throughput cut roughly in half relative to synthetic benchmarks run under ideal conditions. The pattern of the write counts as much as the raw bandwidth available to receive it. Writing thousands of small tensor objects (layer weights, optimizer shards, metadata records) individually generates a flood of small file operations that overwhelms a parallel file system's metadata handling well before it comes close to saturating available data bandwidth.

The storage stack a checkpoint must traverse

Diagram: The Storage Tier Stack: Bandwidth vs. Durability. Visualizes: Visualize the four-tier checkpoint storage hierarchy as a vertical stack, showing how each tier trades bandwidth for durability.

A checkpoint write passes through tiers that differ from one another by orders of magnitude in both bandwidth and latency, and the slowest tier in that chain sets the effective rate for the entire operation, no matter how fast the others run.

GPU HBM sits at the top: fastest, and smallest. Host DRAM comes next, large enough on a single node to hold a full checkpoint copy for a modest-sized model, but it vanishes the instant that node fails. Local NVMe is durable and fast but tied to the node it lives on, useless once that node is unreachable. A shared parallel file system, the Lustre or GPFS or WekaFS or VAST class of system, is durable and accessible cluster-wide, but under concurrent writes from thousands of ranks at once, it becomes the bandwidth bottleneck for the whole path.

The raw numbers at each hop matter. A single PCIe Gen5 NVMe drive peaks around 14 GB/s. Aggregated across a cluster over NVMe-oF fabrics with RDMA interconnects, sustained throughput climbs far past that single-drive figure, as shown by the deployment at NVIDIA's Eos supercomputer, where the storage layer delivers four terabytes per second in aggregate to 576 DGX H100 systems. That latency figure isn't a curiosity: RDMA and NVMe-oF together achieve roughly 10–30 microsecond latencies through CPU-bypass transfers and by avoiding SCSI emulation layers, and per-file metadata overhead multiplies across thousands of concurrent writers, so shaving microseconds off each operation adds up to real wall-clock savings at checkpoint time.

GPUDirect Storage closes the remaining gap by building a zero-copy path straight from NVMe, whether local or reached over NVMe-oF, directly into GPU memory, cutting the CPU out of the staging process entirely on both the write and the restore path. That capability rides on top of a network fabric, not alongside it. NVMe-oF with GPUDirect Storage runs over RDMA transports such as InfiniBand or RoCE, or over standard TCP/Ethernet via NVMe/TCP, and whichever transport a cluster picks is a hard prerequisite for the checkpoint design, not a setting to tune after the fact. Even with asynchronous flush and prefetch in place, the orders-of-magnitude gaps between these tiers still produce bottlenecks under concurrency, and the case where thousands of ranks all try to flush at once, the multi-writer checkpoint storm, remains the failure mode every design in this space has to account for.

Why staging through host DRAM breaks the synchronous barrier

Moving the checkpoint copy off the critical training path and into host DRAM is the single highest-leverage change available to a checkpoint system, and it earns that ranking because it breaks the global synchronous barrier without touching the persistent storage layer at all. Training doesn't have to stall while data crosses a network to a parallel file system. It only has to stall long enough to copy state into local host memory, a transfer where PCIe bandwidth is high relative to the host DRAM write rate, finishing within a small number of training iterations when overlapped with the forward and backward passes already in flight.

This is not a speculative technique. Systems including CheckFreq, DataStates-LLM, and PCcheck all pipeline partitioned checkpoints into host CPU buffers specifically to overlap I/O with ongoing computation, and doing so now counts as the baseline expectation for a production training system rather than an advanced optimization reserved for the largest labs. PHOENIX pushes the idea further by treating the in-memory checkpoint as spatial redundancy that lives entirely off the critical path, overlapping the copy with computation closely enough to reach zero measured overhead during error-free execution, with full recovery completing in under 40 seconds across configurations of up to 512 A100 GPUs and models running up to tens of billions of parameters.

The limitation appears as soon as a node actually dies. Host DRAM does not survive a node failure, so an in-memory checkpoint only protects against transient failures such as process crashes or network resets. A permanent node loss demands that the checkpoint has already made it to durable storage before the node disappears. That constraint creates the defining tension of checkpoint system design: checkpoint often enough to keep recomputation cheap after a failure, but flush to durable storage fast enough to survive a permanent one, and those two goals pull against each other whenever the flush path itself is synchronous. PHOENIX handles the permanent-failure case by pairing its in-memory checkpoint with a communicator reconstruction protocol that hot-swaps a failed node for a spare one at runtime, sidestepping a full job restart. PHOENIX does not close the gap of surviving genuine, permanent node loss with a checkpoint sitting only in volatile memory, and that gap pushes checkpoint designers toward multi-tier storage architectures.

Tiered persistence: aligning storage placement with failure heterogeneity

The right storage tier for any given checkpoint copy depends on which class of failure it needs to survive. Transient process failures, permanent node failures, and rack- or cluster-scale outages each demand a different durability guarantee, and building every tier to withstand the rarest of these failures wastes bandwidth that could otherwise go toward faster, more frequent checkpoints.

TierCheck makes that mapping explicit. It runs a three-tier design that keeps lightweight differential checkpoints in local and peer memory for fast, localized recovery, while asynchronously migrating the heavier base checkpoints to remote persistent storage for durability against permanent loss, all while holding strict global consistency across every tier. In evaluation, that design brought end-to-end checkpointing time under 10 seconds, a figure that reflects what's possible once each tier is asked to do only the job suited to it, rather than every copy paying the cost of the most durable, slowest path available.

SPARe attacks the same problem from a different angle. Instead of optimizing the checkpoint path itself, it masks node failures during gradient synchronization by stacking redundant data shards across parallelism groups and adaptively reordering execution around the loss, cutting time-to-train substantially compared with traditional replication schemes at extreme scale. SPARe treats multi-level checkpointing, spreading I/O across device, host, and storage tiers, as an established technique that lightens the rework burden once a job resumes, separate from the restart-latency issue SPARe itself is built to solve. The two approaches are complementary rather than competing: one shrinks the cost of the checkpoint path, the other reduces how often a failure forces a restart in the first place.

Portability across configurations is a separate problem tiering alone doesn't solve. ByteCheckpoint offers a unified checkpointing system built for large foundation model development, and AdaCheck brings adaptive checkpointing with redundancy utilization to the same space; both address what happens when checkpoint state has to move across different parallelism configurations, a resharding and consistency challenge distinct from tier placement.

Adding tiers adds failure surface inside the checkpoint system itself, and that's a fair objection. Differential checkpoints sitting in peer memory can reflect a more recent training iteration than the base checkpoint already resting on shared storage, so any recovery has to reconstruct one globally consistent state out of a mix of snapshots taken at different tiers and different times, without letting parameter values drift out of sync across ranks. The alternative, forcing every checkpoint through a single synchronous path to durable shared storage, is demonstrably worse once a job reaches tens of thousands of GPUs, and the consistency protocols tiering requires are a solvable engineering problem rather than an open one.

Building the async flush pipeline: coalescing, queue depth, and I/O backend choice

The stretch of the write path between host DRAM and persistent storage is where most checkpoint implementations lose the performance they worked to gain, and nearly all of that loss traces back to a mismatch between the I/O pattern and the storage hardware rather than any hard limit in the hardware itself.

The fix is aggregation, alignment, and coalescing, which a 2026 study found to be the dominant throughput levers, restoring bandwidth and reducing metadata overhead by aggregating small tensor writes into larger operations, and achieving up to 3.9× higher write throughput than DataStates-LLM and 7.6× higher than TorchSnapshot using liburing-based I/O. The reasoning behind that gain is counterintuitive at first: individual checkpoint tensors, whether layer weights, optimizer shards, or metadata records, are small relative to the I/O size that gets the best throughput out of NVMe drives and parallel file systems. Writing each one as its own operation generates a metadata operation for every single tensor, and at scale that flood of metadata saturates a file system's namespace server long before it ever comes close to saturating the storage's actual data bandwidth. The bottleneck sits in the namespace server's request queue rather than the disks.

Backend choice compounds this. The POSIX I/O interface, still the default in a lot of production code, serializes submissions and completions in a way that keeps applications from exploiting the full queue depth modern NVMe hardware supports. liburing's submission and completion ring model allows batched, asynchronous submissions matched to hardware queue depth, handles completions across multiple threads, and reuses persistent buffers instead of allocating fresh memory on every checkpoint interval.

Buffered and direct I/O trade off in ways that aren't obvious from either one's marketing description. Buffered I/O hides the penalty of small writes by aggregating them in the page cache, but that comes with added memory pressure and flush timing that's hard to predict in advance. Direct I/O demands explicit alignment and coalescing inside the application itself, in exchange for deterministic throughput and no double-buffering overhead. For checkpoint writes specifically, direct I/O paired with explicit coalescing is the better path: checkpoint I/O is write-once, and once the writes are coalesced into large enough operations, they saturate NVMe bandwidth on their own without needing the page cache as an intermediary. Pre-allocating a fixed pool of staging buffers and reusing them across every checkpoint interval removes the malloc and free overhead that otherwise accumulates as checkpoint frequency climbs, a detail that matters more the more often a job checkpoints.

Reducing what gets written: differential and layer-wise selective persistence

Every technique so far makes the write faster. A separate and equally valid lever cuts the volume of the write instead, and writing less data can reduce I/O cost more aggressively than pipeline tuning alone manages, though each way of doing that trades some write-time savings against added recovery complexity or a small accuracy risk.

LayerCheck starts from an empirical observation about how models actually change during training: weight updates across transformer layers during post-training are not distributed evenly. Some layers shift substantially at every checkpoint interval, others barely move at all. Persisting only the layers whose updates cross a set threshold spreads writes out over time instead of bursting all of them at once every N steps, which smooths the I/O profile across the run rather than concentrating it into periodic spikes.

Recovering from a LayerCheck checkpoint is not a simple file read. It requires reconstructing a composite state made of layers persisted at different timestamps, matching each layer's most recently saved version with the corresponding optimizer state, and that mismatch introduces a bounded staleness term that holds under standard Adam optimizer assumptions. That's the trade the technique asks a training system to accept: substantially smaller writes, up to a 22.6× reduction in total checkpoint size, and 1.31× faster end-to-end training, with post-restart loss deviating at most 0.54% from the failure-free trajectory, purchased at the cost of a restart procedure that has to reconcile time-skewed layer versions rather than load one clean, uniform snapshot. For a system already carrying async flush and tiered persistence, that's a reasonable exchange. It's not a free one.

Diagram: LayerCheck's Write Reduction vs. Recovery Trade-off. Visualizes: Show the concrete outcome of LayerCheck's selective persistence approach as a before/after or two-column contrast: traditional full checkpointing (all layers written every N…

Sources

  1. PHOENIX: Resilient LLM Training with Hot-Swapping via Zero-Overhead Checkpoint
  2. Understanding LLM Checkpoint/Restore I/O Strategies and Patterns
  3. LayerCheck: Adaptive Layer-wise Checkpointing for Large Language Model Post-training
  4. SPARe: Stacked Parallelism with Adaptive Reordering for Fault-Tolerant LLM Pretraining Systems with 100k+ GPUs
  5. TierCheck: Tiered Checkpointing for Fault Tolerance in Large Language Model Training
  6. TierCheck: Tiered Checkpointing for Fault Tolerance in Large Language Model Training
Filed underSystems Design

More in Systems Design