All essays
TechnicalDEEP DIVEFEB 2026

Distributed Checkpointing Strategies for Large-Scale GPU Training

Distributed checkpointing methods for multi-GPU training including synchronous vs asynchronous approaches, storage backend optimization, and incremental checkpoint compression.

01

SYNCHRONOUS VS ASYNCHRONOUS CHECKPOINTING

Checkpointing is the dominant failure recovery mechanism for long-running distributed training. Synchronous checkpointing pauses training for 30-120 seconds every checkpoint interval to write optimizer state and model weights to persistent storage. At 512 H100 GPUs training Llama 3 70B, each synchronous checkpoint saves 120 GB and costs $23-$90 in lost training time. Async checkpointing overlaps write with computation, reducing visible pause to 5-15 seconds.

Async checkpointing risks: if training crashes during a checkpoint write cycle, the GPU memory state may be inconsistent with storage. Two-buffer checkpointing mitigates this by maintaining a committed checkpoint and building the next in a separate buffer. Committed checkpoints are always consistent. For Llama 3 70B with two-buffer async, maximum uncommitted work is 30 seconds versus 15 minutes for synchronous.

Checkpoint MethodTraining PauseMax Data LossStorage IOPS RequiredBest Cluster Size
Synchronous30-120 sec15 min20,000-40,0008-128 GPUs
Async (single buffer)5-15 sec15 min + write40,000-80,00064-512 GPUs
Async (two-buffer)5-15 secCheckpoint interval40,000-80,000128-2,048 GPUs
Async incremental2-5 secCheckpoint interval10,000-20,000512-8,192 GPUs
Async + compression3-10 secCheckpoint interval5,000-15,000256-4,096 GPUs
02

STORAGE BACKEND OPTIMIZATION

Checkpoint storage backend choice significantly impacts performance. Parallel filesystems like Lustre and GPFS provide 50-100 GB/s write throughput but cost $2-$5 per GB-month. Object stores like S3 provide 5-15 GB/s write throughput at $0.023 per GB-month but add 200-800ms per PUT operation latency. Hybrid tiering stores the latest checkpoint on parallel filesystem with older versions in object storage.

NVIDIA GPUDirect Storage enables direct GPU-to-storage data transfer bypassing CPU memory. GDS reduces checkpoint write time by 30-50 percent for large models. Combined with asynchronous checkpointing, GDS achieves checkpoint completion in 3-8 seconds versus 45-90 seconds without GDS. Implementation requires compatible filesystem (Lustre 2.14+, GPFS 5.1+) and storage over NVMe or InfiniBand.

03

CHECKPOINT COMPRESSION AND ENCODING

Checkpoint size reduction through compression saves storage and bandwidth. zstd at compression level 3 achieves 2.5-3.5x compression on FP16 optimizer states, reducing 120 GB to 35-48 GB. The compression overhead of 5-15 seconds is offset by 60-70 percent faster transfer to object storage. Delta encoding storing only parameter changes since last checkpoint achieves 80-90 percent reduction.

Lossy checkpoint compression using FP8 storage for optimizer states achieves 2x additional compression with negligible training impact. Mixed-precision checkpointing stores master weights in FP32 but optimizer momentum in FP8, saving 8 bytes per parameter. For Llama 3 70B with 70 billion parameters, this reduces checkpoint size from 120 GB to 75 GB.

04

FAULT TOLERANCE AND RECOVERY AUTOMATION

Automated recovery from checkpoint reduces MTTR from hours to minutes. Elastic training frameworks like PyTorch FSDP with torchtitan support automatic checkpoint loading and training resumption upon cluster failure. Recovery time at 256 GPUs is 2-5 minutes versus 15-45 minutes for manual recovery. Checkpoint validation via SHA-256 hash verification catches 99.97 percent of storage corruption incidents.

The checkpoint frequency optimization problem balances compute loss vs storage cost. Given 2 percent chance of failure per hour, 15-minute checkpoint frequency yields 0.5 percent compute loss from checkpoint overhead and 0.3 percent from failure recovery, totaling 0.8 percent overhead. Reducing to 60-minute frequency cuts checkpoint overhead to 0.12 percent but increases failure loss to 1.2 percent, totaling 1.32 percent.

Filed under
CheckpointingDistributed TrainingPyTorch FSDPAsync CheckpointModel ParallelismFault Tolerance