Chapter 9.4
In this chapter · 8 sections
Checkpointing for Large-Scale Training
Checkpointing is a training cluster's goodput control knob — the interval, tier, and write bandwidth you choose decide how many GPU-hours each failure erases, and at frontier scale failures arrive constantly.
What you'll decide here
- How much checkpoint state you actually carry — about 14 bytes/parameter for the BF16 weights plus FP32 master weights and Adam moments recipe derived here — because those unique bytes, save frequency and drain deadline set the checkpoint-path bandwidth and the failure domains its committed copies must survive.
- The synchronous checkpoint interval: use Young's first-order approximation √(2·C·MTBF) with C measured through durable completion, or Daly's higher-order model when failures during save and restart matter; evaluate asynchronous saving against its separate captured-state-age, drain, commit and restart limits.
- The tiering architecture — synchronous vs asynchronous vs in-memory/peer, and how many storage tiers — that determines whether a failure costs you minutes or tens of minutes of replay.
- The on-disk format: monolithic single-rank gather vs sharded/distributed (per-rank) checkpoints, and whether the format survives a reshard to a different parallelism layout on restart.
- Whether compression or sparsification is worth the CPU and the fidelity risk, versus simply provisioning more drain bandwidth — compare lossless compression, incremental and selective saving against total recovery cost, and test any lossy optimizer-state or MoE-specific scheme for restart fidelity.
At small scale, checkpointing is an afterthought: save the model every so often, restore if something breaks. At frontier scale it becomes one of the load-bearing decisions in the whole training system, because the arithmetic flips. A 16,384-GPU run failed roughly once every three hours in Meta's Llama 3 405B snapshot — 419 unplanned interruptions over 54 days, 78% of them hardware-caused (Chapter 10.7). Larger synchronous jobs expose more components and shared dependencies, but their interruption cadence must be measured for the named fleet, job membership, software stack, and event definition rather than inferred from accelerator count alone. A fault that escapes local retry, spare substitution or the runtime’s fault-tolerance mechanism can become a full-job stall for a synchronous run: the whole machine rewinds to the last good checkpoint and replays the work in between. Count those job interruptions, not every component alarm. The question is no longer whether to checkpoint but how often, to where, and at what bandwidth — and getting those three wrong directly reduces measured goodput.
This chapter is the canonical home for the checkpoint math the rest of the guide cross-references, built decision by decision — the anatomy of a checkpoint (why ~14 bytes/parameter, and what is in those bytes), the failure model that makes checkpointing a goodput problem rather than a durability one, Young's first-order approximation and Daly's higher-order checkpoint model that turns MTBF and save-cost into a frequency, the tiering spectrum from synchronous-to-object up through in-memory peer copies, the format fork between monolithic and sharded checkpoints, compression and sparsification, and finally how to size the bandwidth and isolate the incast it creates.
Checkpoint anatomy: why ~14 bytes per parameter
The first sizing number you need is how big a checkpoint is, and the starting point is a declared serialization recipe: ~14 bytes of checkpoint state per model parameter — the serialized-state recipe VAST Data assumed in order to infer model sizes from a survey of 85k+ checkpoints, carrying about ±15% uncertainty, and the one to verify against what your framework actually writes. It is not the model weights — those are the smallest part. The breakdown, for the standard BF16-compute / FP32-master / Adam-optimizer recipe, is:
- 2 bytes — the BF16 (or FP16) model weights used in the forward/backward pass.
- 4 bytes — the FP32 master copy of the weights the optimizer actually updates (mixed precision keeps a high-precision master to avoid drift).
- 4 bytes — Adam's first moment (the running mean of gradients), one FP32 value per parameter.
- 4 bytes — Adam's second moment (the exponential moving average of squared gradients), again FP32 per parameter.
That sums to 14. Gradients are usually not checkpointed — they are recomputed on restart from the restored weights and a forward/backward pass, so they cost replay time, not storage. The implication that catches teams out: Under this recipe, a checkpoint is 7× the BF16 weight file you think you are saving. A 405B-parameter model is a ~0.8 TB weight file but a ~5.7 TB checkpoint; a 1T-parameter model is a ~14 TB checkpoint. Optimizer state, not weights, dominates the I/O — which is why "just save the weights" is the wrong mental model for resumable training, and why optimizer-state sharding (ZeRO/FSDP) is as much a checkpoint-I/O strategy as a memory strategy.
The decision this forces is upstream of storage: if you change the optimizer or the precision recipe, you change the checkpoint size and therefore the bandwidth budget. FP32 Adam is 4 + 4 bytes for the moments, plus a 4-byte master weight and a 2-byte BF16/FP16 weight: 14 bytes per parameter. 8-bit Adam is 1 + 1 bytes for the moments, plus the same 4-byte master and 2-byte weight: 8 bytes per parameter; a memory-efficient optimizer that drops the second moment changes the arithmetic again. These are training-science choices with a direct, often-unnoticed consequence on the storage and fabric you must provision. The parameter-to-bytes mapping is the bridge between the two.
The failure model: checkpointing is a goodput problem
To size a checkpoint interval you need a failure model, and the one that matters for synchronous training is starkly simple: a fault that escapes local recovery can force whole-job rollback. One GPU falling off the bus, one optical transceiver degrading, one HBM uncorrectable error — when redundant paths and local recovery cannot absorb the fault, a synchronous job spanning tens of thousands of GPUs can roll back to its last checkpoint and replay. This is a defining property of synchronous training and the reason checkpoint interval matters: it controls the lost-work component of a job interruption. It does not select the facility topology; N, N+1, distributed, or 2N still follows from the permitted maintenance and fault states, interruption and recovery limits, post-event capacity, independence, common modes, and contract.
The cost of a single failure, in lost GPU-hours, has three parts: the wasted compute since the last checkpoint (on average half the interval), the detection-and-restart latency (the time to notice the fault, evict the bad node, and reload state), and the checkpoint overhead paid during normal operation, accounted separately from the incident duration. Push the interval long and you save overhead but risk losing a lot of work per failure; push it short and you lose little work but pay overhead constantly. The optimum is where these balance, and it changes with the effective interruption distribution measured for the actual job. Published anchors from different fleets are not one arithmetic curve: correlated faults, shared services, software, detection policy, and job membership all change job-level MTBF. Re-estimate that input from telemetry as the fleet and workload change.
Scope & caveats
An arithmetic recipe for the standard BF16-compute / FP32-master / Adam mixed-precision configuration, not a measured population value. VAST Data assumed 14 bytes/param to infer model sizes from checkpoint sizes (its own footnote puts the resulting uncertainty at about ±15%); the survey did not establish the recipe. 8-bit Adam moments give 8 bytes/param; optimizers that drop the second moment change it again. Verify the tensors the chosen framework actually serializes.
Scope & caveats
SemiAnalysis reported this for one 512-H100 cluster at a top-tier operator in October 2024. It is not a per-GPU rate, a scaling law, or a portable fleet baseline; do not extrapolate it by accelerator count.
Scope & caveats
Define the project drain window from snapshot release through background drain completion, reserve the durable commit tail separately, and add measured contention and concurrency headroom. Do not substitute 10% of the checkpoint interval.
Scope & caveats
Google reports a 6.59% increase in ML Goodput on one 35K-chip TPU v5p workload; that wording does not establish 6.59 percentage points absolute, and the result is specific to that workload, checkpoint cadence and storage tier. Vendor-reported, not independently measured.
Scope & caveats
Stipulated endpoints for sensitivity only. They are neither measured provider outcomes nor universal targets. Chapter 14.1 reconciles productive-time boundaries; a site measures its own baseline.
Young first-order and Daly higher-order checkpoint intervals
The interval is not a matter of taste. Young's 1974 first-order approximation gives the square-root interval below; Daly's 2006 higher-order model uses a fuller checkpoint/restart cost model. This is the canonical result this guide refers back to from Chapters 1.2, 10.7, and 14.4 — and it falls straight out of the failure model above. The intuition: the wasted-work cost per unit time grows with the interval (longer interval, more to replay on failure), while the checkpoint-overhead cost per unit time shrinks with the interval (fewer, more widely-spaced saves). Minimize the sum and you get a square-root law:
τopt ≈ √( 2 · C · MTBF )
where τopt is useful compute time between synchronous saves, C is blocking save time through durable completion, and M = MTBF is whole-job mean time to interruption. The failure-free wall-clock cycle is τ + C; its save fraction is C / (τ + C). The optimal interval scales with the square root of both the save cost and the MTBF. Halve your checkpoint write time (faster tier, more drain bandwidth) and the optimal interval shrinks only by √2 — and the optimal overhead fraction, √(2C/M), falls by the same √2, so halving C buys a √2 improvement, not a twofold one. If individual accelerator and node hazards were independent and identically distributed and membership were the only scaling variable, doubling cluster size would approximately halve whole-job MTBF and shrink the square-root interval by √2. Real jobs also fail through correlated hardware, shared services, software, membership changes, and detection policy, so size the interval from the named fleet's measured whole-job MTBF rather than scaling a component or smaller-cluster observation.
Either interval model is only as good as the measured whole-job failure and recovery inputs you feed it, and most teams feed it a fantasy. A cluster whose burn-in telemetry shows elevated interruptions needs a different cadence; a run that hits a lemon node (Chapter 10.7) sees its effective MTBF crater. The right practice is to measure MTBF continuously from the fleet's own fault telemetry and re-derive τ as conditions change — early in a run, interval shorter; once burn-in clears, interval can lengthen. A static interval set at scoping time is wrong on day one and wrong again at steady state.
For the synchronous Poisson-interruption model, Daly (2006), equation 37 gives τ ≈ √(2CM)[1 + √(C/(2M))/3 + C/(18M)] − C when C < 2M, and τ ≈ M otherwise. Failures during restart are included: restart duration changes total completion time but cancels from this optimum. Asynchronous capture needs a separate protected-state timeline; replacing C by a short local snapshot does not optimize durable recovery.
Worked decision: interval, captured progress and protected commitment
The synchronous comparison blocks through capture, drain and commit: C = 20 + 60 + 5 = 85 s. Young gives √(2 × 85 s × 1,800 s) ≈ 550 s of useful compute, with a roughly 640 s cycle. That unconstrained candidate already exceeds the 360 s protected-state-age limit before adding the next save's commitment lag. A synchronous policy constrained to the same captured-state bound requires S ≤ 360 − 85 = 275 s, but its save fraction is then at least 85/275 ≈ 31%, above the 10% allowance. Faster durable saving or a different policy is required.
Choose the assumed asynchronous policy at S = 270 s. Unique bytes are 100 × 109 parameters × 14 B/parameter = 1.4 TB, partitioned across writers. Drain needs 1,400 GB / 60 s ≈ 23 GB/s of protected logical payload. Commitment occurs 20 + 60 + 5 = 85 s after capture, leaving 270 − 85 = 185 s before the next capture. Immediately before the next commit, the last committed captured state is S + Ccapture + D + K = 355 s old, within 360 s. Its stored step and consumed-sample cursor identify the actual progress recovered; wall-clock age is a conservative progress-loss bound that includes time spent paused.
The capture fraction is 20/270 ≈ 7.4%, before failures and contention; asynchronous drain earns the choice by fitting the pause budget. Two committed generations plus one in flight require 3 × 1.4 = 4.2 TB logical before protection and reserve. Restore requires its own 1,400/60 ≈ 23 GB/s read path; 30 + 30 + 60 + 30 = 150 s meets the 180 s RTO. These calculations pass the assumed policy, while hardware acceptance remains HOLD until the full timeline and state restore are measured.
Flip: lengthening S to 300 s gives 300 + 85 = 385 s of state age and fails. The crossover is S = 275 s; clearing every drain is insufficient by itself. At S = 270 s, drain must stay at or below 65 s to preserve the age limit. Keep S = 270 s and qualify the 60 s drain budget; if the measured tail misses it, shorten the period only after checking pause and backlog limits, or buy more service capacity. The interval method is Young/Daly's synchronous model; this chapter owns the asynchronous age ledger, with orchestration handed to Chapter 10.7.
Synchronous vs asynchronous vs in-memory: the tiering spectrum
The biggest lever on checkpoint cost C is where the checkpoint lands and how much of the write blocks training. There is a clear spectrum, and most large runs use several tiers at once.
Synchronous checkpointing stops training, writes the full state to durable storage, and resumes. It is simple and it is correct, but it is the most expensive: the GPUs sit idle for the entire write, so C is the full write time and it grows with model size. At frontier scale a synchronous-to-object checkpoint can stall the machine for tens of minutes — unacceptable when MTBF is also tens of minutes. Synchronous-only checkpointing remains viable even on a large cluster when the complete save fits its pause and recovery budgets.
Asynchronous checkpointing has four distinct timing terms. The blocking snapshot copies GPU state to host RAM or local NVMe and pauses training only for that capture. Training then resumes while the background drain moves the snapshot to durable storage; define an explicit drain-completion window D from snapshot release to the deadline for background drain completion. The separate durable commit tail covers flush, metadata publication, verification, and atomic visibility before the checkpoint is recoverable. Finally, add contention and concurrency headroom for simultaneous checkpoints and shared fabric or storage traffic. Size the ideal drain floor as unique serialized bytes divided by D. Because D starts at snapshot release and ends at drain completion, capture and commit are already outside it; subtract them only when converting a capture-start-to-commit budget into D. Add contention and concurrency headroom separately. VAST reported 50–200 GB/s drain rates across 40 production runs and a 10% overlap observation; those are reported workload results, not a universal fraction of the checkpoint interval. Asynchronous checkpointing lets training continue while the captured state drains; choose it when protected commit and interference fit the budgets better than a qualified synchronous save.
In-memory / peer checkpointing goes further: it keeps the most recent checkpoint in the CPU DRAM of peer nodes rather than (or in addition to) durable storage. Because aggregate host-memory bandwidth across the cluster dwarfs any storage tier, the checkpoint can be taken every single step with near-zero throughput overhead (the GEMINI/in-memory line of work demonstrated optimal per-iteration checkpointing and >13x faster recovery). On a failure, a surviving peer holds a copy of the dead node's shard, so recovery is a memory-to-memory restore in seconds rather than a storage read. The cost is RAM: you reserve host memory across the fleet to hold redundant copies, and you need a placement strategy that guarantees the failure of any node leaves its shard recoverable from a survivor. This is the frontier-operator pattern — capture intervals approaching the step duration, with surviving peer placement defining how much progress remains recoverable.
| Tier | Write target | Blocking cost C | Recovery latency | Best for |
|---|---|---|---|---|
| Synchronous to durable | Parallel FS / object store | Full write (minutes at scale) | Full read from durable | Small clusters; final / milestone checkpoints |
| Async drain to durable | Local snapshot → background drain to FS/object | Seconds (snapshot only) | Read from durable (minutes) | The default working tier for large runs |
| Local NVMe scratch | Node-local NVMe (per-rank shard) | Seconds (local write) | Local re-read if node survives; else next tier | Fast restart when the fault is recoverable in place |
| In-memory / peer copy | Peer-node host DRAM (redundant shard) | Near-zero (every-step capable) | Memory-to-memory (seconds) | Hyperscale runs where MTBF approaches step time |
Sharded vs monolithic: the format fork
Independent of where the checkpoint lands is the question of how it is laid out, and this is a genuine fork with a sharp scaling consequence. A monolithic checkpoint gathers the full model and optimizer state onto one rank (historically rank 0), which serializes a single large file. It is easy to reason about and easy to load anywhere, but the all-gather is a bottleneck that does not scale: as the cluster grows, more ranks funnel state through one writer, and C grows with model size while the write bandwidth stays pinned to a single node's link. Monolithic checkpoints are a small-model convenience that breaks down when gathering state into one writer exhausts its memory or misses the checkpoint write budget.
A sharded / distributed checkpoint writes each rank's own shard in parallel — every GPU drains its slice of weights and optimizer state to storage simultaneously. Aggregate write bandwidth now scales with the cluster, which is the only way to keep C small at frontier scale. This is the format behind PyTorch Distributed Checkpoint (DCP), DeepSpeed/Megatron sharded checkpoints, and the Orbax format on the JAX/TPU side. The cost is complexity: a sharded checkpoint is a directory of many files plus metadata describing the parallelism layout that produced it, and reading it back is a distributed operation, not a single load.
The consequence that bites teams hardest is resharding. A sharded checkpoint written under one parallelism layout — say TP=8, PP=16, DP=64 — is not trivially loadable under a different layout. If a failure forces you to restart on fewer nodes, or you want to resume a run with a different tensor/pipeline split, a naive sharded format requires a slow offline reshard or cannot load at all. Modern distributed-checkpoint libraries solve this with a layout-independent logical format: the checkpoint records the global tensor shapes and each shard's coordinates, so the loader can re-partition on the fly to the layouts supported by that library, state representation and model implementation. Verify your checkpoint format survives a reshard before you need it to — discovering at 3 a.m. that your only checkpoint cannot load on the 12,000 GPUs you have left (instead of the 16,000 you started with) is a goodput catastrophe that elastic, reshard-capable formats exist specifically to prevent.
Deep dive: the metadata and consistency trap in sharded checkpoints
A sharded checkpoint is valid only when every required shard and its application state belong to the same captured step. Imagine rank 17 finishing its drain in 4 seconds, rank 4,002 taking 9, and a failure landing between them: publishing that directory as complete would expose a torn checkpoint. Freeze a generation manifest containing tensor names, global shapes, dtypes and shard coordinates; include optimizer, scheduler, random-number and mixed-precision scaler state, dataset version and consumed-sample cursor. Serialize each unique tensor once; record intentional replicas separately. Writer count controls placement and parallelism, not model size. The PyTorch DCP save/load contract covers distributed state; the application owns the state it passes and the versions/layouts it can restore.
Write shards to immutable generation-specific locations, verify their checksums, then publish the complete manifest. On a filesystem, qualify data and directory persistence and the supported atomic publication operation. On an object store, complete all uploads before publishing a manifest; use conditional publication to prevent concurrent writers from replacing the chosen generation. An object prefix is not an atomic directory rename. S3 conditional writes supply overwrite controls, not a transaction across all shard objects.
Recovery selects the newest complete, verified generation whose placement survives the event. A partial newer generation never supersedes the prior committed one. Retain an older valid generation for corruption or failed validation and enough space for the next save; the retention count is a policy input. Publication can finish in the background, so count a commit tail in training pause only if the runtime actually blocks on it. Inject missing shards, a writer death before publication, corrupt data after publication and changed restart membership. Require rejection of the damaged generation and an exact next-sample/step validation before declaring recovery successful. State age, recovery time and application correctness are separate pass gates.
Compression, sparsification, and when they pay
Because the checkpoint is dominated by optimizer state, an obvious lever is to make it smaller — compress it or save only part of it. For most runs the better move is to provision more drain bandwidth before you reach for compression, because lossless compression of BF16/FP32 tensors yields modest ratios at real CPU cost, and that CPU competes with the data loader (Chapter 9.5) for the same host cores. Halving the checkpoint size halves C, and the square-root structure of Young's first-order approximation then moves the optimal interval by only √2 — a smaller win than the bandwidth and CPU spend often justify.
There are specific cases where reducing checkpoint volume does pay. Optimizer-state quantization (8-bit Adam moments) cuts the two moment tensors from 4 + 4 to 1 + 1 bytes per parameter, reducing the full recipe from 14 bytes to 8 while retaining the 4-byte FP32 master and 2-byte BF16/FP16 weight — a training-science decision with a storage payoff. Sparse / incremental checkpointing snapshots only the subset of state that changed (or a rotating subset of operators) each interval rather than the full state every time, which is especially attractive for Mixture-of-Experts models. Two different mechanisms hide under that heading and they carry different recovery contracts. Incremental checkpointing writes state that changed since its valid base; reconstruction needs the base and every required delta. Keep that dependency chain protected and verify reconstructed state before claiming unchanged restart fidelity. Selective expert checkpointing — the MoEtion line of work — exploits expert popularity skew to snapshot hot experts often and cold ones rarely, buying high frequency by tolerating a bounded amount of dropped work on the cold path. Over a full batch most experts do see updates, so the second mechanism is not free: decide explicitly what may be lost on restart and which convergence test justifies losing it. The split is clean: for dense models, buy bandwidth; for MoE and 8-bit-optimizer regimes, the structure of the state itself makes reducing checkpoint volume the better lever.
Restart, recovery, and sizing the bandwidth
Writing checkpoints is half the system; recovery is the half that determines how long a failure actually costs. A complete recovery is: detect the fault, isolate and evict the bad node (ideally swapping in a hot spare so the parallelism layout is preserved), reload state to every rank, recompute the gradients the checkpoint did not store, and resume. Detection latency is its own discipline (Chapter 10.7) — silent data corruption and slow-degrading hardware can corrupt steps for a long time before a checkpoint is even triggered. The reload itself is where checkpoint architecture pays off: a peer/in-memory restore is seconds, a local-NVMe re-read is seconds-to-a-minute, and a cold read of a multi-terabyte sharded checkpoint from the durable tier is minutes — the same 30-minutes-to-under-a-minute spread the tiering decision controls.
Define the background window from snapshot release to drain completion, as in the worked policy. Convert a whole capture-to-commit deadline by removing capture and commit once, then divide unique payload by the remaining drain time. Concurrent jobs share the delivered write budget, so sum the overlapping streams and include measured contention. Keep read bandwidth for recovery separate: ranks reload concurrently, but the traffic can be storage fan-out, client incast or both depending on shard placement. Size the actual bottleneck cuts and the reload deadline. A write-optimized tier whose readers miss that deadline still fails recovery even if its last checkpoint was durable.
And that burst is the part that breaks naively-designed clusters. A synchronized checkpoint write creates incast at shared storage destinations; recovery reads create fan-out and receiving-side concentration that, if it shares the back-end fabric with the training collectives, will collide with the all-reduce traffic and tank the very goodput it exists to protect. The co-design rule is to isolate checkpoint I/O from the compute fabric — a dedicated storage rail, or strict QoS separation on a converged fabric — so a checkpoint storm never starves an all-gather. This is the fabric-and-storage co-design problem treated in Chapter 8.5 and revisited for storage sizing in Chapter 9.8; the point here is that the checkpoint interval you choose creates a periodic incast whose magnitude you computed above, and the fabric must be designed to absorb it without disturbing training.
Putting it together: the checkpoint policy as a goodput contract
A defensible checkpoint policy is a coupled set of decisions, not a single setting — together they define how much of a failure is recoverable. Size the unique state: about 14 bytes/parameter for BF16 weights plus FP32 master weights and Adam moments, then add the application state actually saved. Measure MTBF from live fault telemetry and recalculate the interval with the selected Young first-order or Daly higher-order model as burn-in clears. Make C small where asynchronous capture and drain reduce measured training pause, use a sharded, reshard-capable format, and add an in-memory/peer tier when its surviving placement and timed reload let the common single-node failure recover inside the budget; the protected-state timeline still governs the durable copy. Choose the durable-copy cadence from the required failure-domain state age, with a commit protocol so you never load a torn checkpoint. Then size both the drain bandwidth and the recovery incast, and isolate them from the compute fabric. These mechanisms can materially improve measured goodput; Google's named 35K-chip TPU v5p workload reported a 6.59% increase.
Choose the policy whose captured progress, protected commit, pause fraction and recovery timeline all pass the named failure. A short local snapshot buys useful compute only while drain and publication keep the recoverable generation current; otherwise it buys an expanding recovery gap. Hand the generation manifest and timing limits to Chapter 10.7 for orchestration, Chapter 12.3 for regional recovery, Chapter 13.9 for acceptance and Chapter 14.4 for incident execution.
Scope & caveats
Illustrative inputs; qualify project rates and limits before selection.
Cite this chapter
Fehn, J. (2026). Checkpointing for Large-Scale Training (Chapter 9.4). The Definitive Guide to AI Data Centers. https://aidatacenterguide.com/part-9-storage-and-data/9-4-checkpointing-for-large-scale-training (accessed 2026-09-29).
@misc{aidc-9-4,
author = {Fehn, Jacob},
title = {Checkpointing for Large-Scale Training (Chapter 9.4)},
howpublished = {The Definitive Guide to AI Data Centers},
year = {2026},
url = {https://aidatacenterguide.com/part-9-storage-and-data/9-4-checkpointing-for-large-scale-training},
note = {Accessed 2026-09-29}
}