The Checkpointing Problem at Scale
Long-running AI training jobs -- those exceeding 24 hours -- face near-certain failure from at least one hardware event. A 1,024-GPU cluster training a large model for 30 days has a 99.7% probability of experiencing at least one GPU ECC error, node failure, or network partition during the run. Without checkpointing, a failure after 29 days loses 29 days of compute -- approximately $1.2M in GPU time alone.
Checkpointing solves this by periodically saving the model state (weights, optimizer state, random number generator state) to persistent storage. On failure, training resumes from the most recent checkpoint, losing only the compute time between checkpoints. The art is selecting the checkpoint frequency that minimises the expected cost of lost compute versus the overhead of writing checkpoints.
This post covers checkpointing architectures, storage infrastructure, incremental techniques, and recovery strategies for large-scale GPU training at mid-2026.
Checkpoint Frequency and the Cost Trade-Off
Checkpointing overhead directly reduces training throughput. A full model checkpoint for a 70B-parameter model in BF16 is 140 GB. Writing 140 GB to a parallel filesystem with 20 GB/s write bandwidth (typical for Lustre) takes 7 seconds. On a 1,024-GPU cluster training at 50% model FLOPS utilisation, those 7 seconds represent approximately $800 in GPU time for each checkpoint.
The checkpoint frequency decision balances two costs: the cost of checkpointing overhead versus the expected cost of lost compute from failures. The formula is: optimal interval = sqrt((2 * expected compute lost per failure * failure rate) / checkpoint cost per second). For a 1,024-GPU H100 cluster with a mean time between failures of 72 hours, checkpoint cost of $800 per write, and average failure recovery loss of 50% of the checkpoint interval, the optimal checkpoint frequency is approximately every 4-6 hours.
Real-world practice at mid-2026 varies by application. Teams training models requiring weeks of continuous compute tend toward 2-4 hour checkpoints for safety. Teams with reliable hardware and fast recovery pipelines use 8-12 hour intervals. The trend is toward adaptive checkpointing that adjusts frequency based on observed failure patterns.
Incremental and Asynchronous Checkpointing
Full model checkpoints are expensive. Incremental checkpointing saves only the changed portions of model weights since the last full checkpoint. For large models with optimizer states that change slowly (particularly at later training stages), incremental checkpoints can be 10-30% of the full checkpoint size, reducing write time and storage overhead proportionally.
Asynchronous checkpointing overlaps checkpoint writes with training computation using double-buffering: while the GPU is computing the next training step, the previous step's model state is being copied to CPU memory and then to persistent storage. This eliminates the synchronous write stall, making checkpoint overhead negligible (0.5-2% throughput impact) compared to synchronous checkpointing (3-8% impact).
At mid-2026, asynchronous checkpointing is supported by PyTorch Distributed (`torch.distributed.checkpoint`) with async path, NVIDIA NeMo, and DeepSpeed. The key infrastructure requirement is sufficient CPU memory bandwidth to receive GPU-to-CPU transfers without blocking the training pipeline -- typically requiring 256-512 GB of system RAM per GPU node for large models.
Storage Architecture for Checkpoints
Checkpoint storage requires high write throughput, high reliability, and sufficient capacity. A single training run on 1,024 GPUs over 30 days with 4-hour checkpoint intervals produces 180 checkpoints at 140 GB each = 25 TB of checkpoint data. With multiple concurrent training runs, total checkpoint storage reaches 100-500 TB.
The storage architecture recommendation at mid-2026 is a parallel filesystem (WekaFS, Lustre, or GPUDirect-compatible storage) with at least 20 GB/s write bandwidth per 1,000 GPUs, backed by NVMe flash storage with 5-10 TB usable capacity per training node, and object storage (S3-compatible) for checkpoint archival and long-term retention.
The table below compares checkpoint storage options. Most production clusters use a tiered approach: parallel filesystem for active training checkpoints, local NVMe for hot checkpoint cache, and object storage for archival.
| Storage Tier | Write Bandwidth | Latency | Cost/GB/Month | Use Case |
|---|---|---|---|---|
| Local NVMe (per node) | 6-14 GB/s | <10 us | $0.10-0.15 | Checkpoint staging |
| Parallel filesystem (GPUDirect) | 20-100 GB/s | <1 ms | $0.05-0.08 | Active checkpoint storage |
| NFS over RoCE | 2-10 GB/s | 1-5 ms | $0.03-0.05 | Model artifact storage |
| Object storage (S3) | 1-5 GB/s per prefix | 10-100 ms | $0.01-0.03 | Checkpoint archival |
Recovery Strategies: Minimizing Time-to-Resume
When a training job fails, the objective is to minimise time-to-resume (TTR): the wall-clock time from failure detection to the first post-recovery training step. TTR comprises: failure detection time (seconds to minutes), checkpoint identification (selecting the most recent valid checkpoint), checkpoint loading (reading from storage to GPU memory, 10-60 seconds for 140 GB), validation (verifying checkpoint integrity, 5-30 seconds), and training resumption (rebuilding NCCL rings and optimizer state, 30-120 seconds).
The total TTR target should be under 5 minutes. Beyond 5 minutes, the cost of recovery time starts to approach the cost of the checkpoint interval itself, suggesting the checkpoint frequency should be increased. Teams achieve sub-5-minute TTR through: checkpoint integrity validation at write time (not load time), pre-loaded checkpoint cache on each node's local NVMe, GPU memory pre-allocation for the checkpoint buffer, and NCCL ring topology caching across restart cycles.
Automated failure recovery is the norm at mid-2026. Kubernetes-based training operators (Kubeflow, Volcano, or Ray Train) detect job failures, identify the latest checkpoint from the registry, and relaunch the training job with the checkpoint path as the starting point. The entire cycle from failure detection to training resumption should be sub-2-minutes in a well-automated environment.
Checkpoint Validation and Integrity
A corrupted checkpoint is worse than no checkpoint at all because it introduces silent data corruption into the model. Checkpoint corruption can occur from storage hardware errors, network packet corruption during write, or GPU memory errors that are written into the checkpoint before detection.
The integrity strategy includes: CRC checksum on each checkpoint shard at write time, verified on read; model weight convergence check (load checkpoint, run 10 inference steps, compare logits to expected distribution); and replica checkpoint writing to two independent storage paths (e.g., parallel filesystem + object storage). The dual-write approach doubles checkpoint time but provides a fallback if one storage path is corrupted.
At mid-2026, most checkpoint corruption events (approximately 70%) are detected within the first training step after recovery, as the optimizer state divergence triggers NaN loss values. However, the 30% that go undetected for multiple steps can cause significant training degradation. The recommended practice is to run a 50-step validation training cycle after each checkpoint load, comparing loss trajectory to the pre-failure trajectory before declaring the checkpoint valid.
Multi-Node and Elastic Checkpointing
As training scales beyond a single node, checkpoint coordination becomes more complex. Distributed checkpointing requires all GPUs to reach a consistent state simultaneously. The synchronous barrier approach pauses all GPUs, collects the checkpoint from each GPU's shard, and writes them to storage coherently. This is simple but has poor scaling characteristics -- at 1,024 GPUs, the synchronisation barrier alone can add 5-15 seconds of idle time.
Elastic checkpointing supports training jobs that can scale up or down based on GPU availability. This is particularly relevant for spot instance-based training, where GPUs may be preempted mid-run. Elastic checkpoint saves the model state in a way that can be loaded on a different number of GPUs for the next training cycle. The checkpoint contains the full model weights (independent of the number of GPUs) rather than sharded weights.
Elastic checkpoint support in PyTorch (`torch.distributed.elastic`) and DeepSpeed enables training jobs on spot GPU instances with near-100% completion rates despite frequent preemption. The trade-off is that elastic checkpoints are larger (each GPU saves the full model rather than a shard) and require more storage bandwidth during writes. For teams committed to spot-based training, this trade-off is acceptable given the 60-80% cost savings versus reserved pricing.
