SYSTEM Cited by 1 source
PyTorch Distributed Checkpoint (DCP)¶
PyTorch Distributed Checkpoint (DCP) is the PyTorch API for saving and loading
training state as a sharded, parallel checkpoint: every rank writes its own
shard concurrently, alongside a small .metadata file describing how the shards
compose into full tensors. It is the mechanism behind the
distributed checkpoint concept, and the
recommended replacement for monolithic torch.save-on-rank-0 in
PyTorch training. (Source:
sources/2026-08-28-databricks-fast-fault-tolerant-pytorch-training-on-ai-runtime)
What it does¶
- Parallel shard writes. Save time drops roughly as 1/N with rank count,
versus
torch.save's serial gather-to-rank-0-then-write-one-file. - Reshardable load. Because
.metadatarecords the global tensor layout, a DCP checkpoint can be reloaded onto a different number of GPUs — DCP re-plans which bytes each new rank needs, so recovery onto a reduced-capacity cluster after node loss works without manual surgery. async_save. Splits a save into a fast staging copy plus a background upload that overlaps continued training, so the training loop pays only for staging — see asynchronous-checkpoint-staging.- Completeness marker.
.metadatais written only after all shards land, so its presence is a reliable "save complete" signal for automatic recovery.
Applies even to DDP¶
DCP is not just for sharded models. Even a data-parallel (DDP) job — where every rank holds identical weights — benefits: DCP shards the model state and writes it in parallel anyway. Critically, it's the same API used for FSDP and tensor parallelism, so adopting it on a DDP job means never rewriting resilience code when the model later shards.
On Databricks AI Runtime¶
AI Runtime implements DCP against Unity Catalog
volumes via UCVolumeWriter and UCVolumeReader, which stage I/O through
local NVMe and mark a checkpoint complete only once its data has fully landed.
Measured async_save vs torch.save (excluding the baseline's network-storage
time):
| Job | Speedup |
|---|---|
| DDP LLM, 2.8B params on 32×H100 | 1.8× (36s vs 66s) |
| FSDP LLM, 20B params on 32×H100 | 58× (522s vs 9s) |
Seen in¶
- sources/2026-08-28-databricks-fast-fault-tolerant-pytorch-training-on-ai-runtime
— DCP vs
torch.save,async_savestaging,.metadataresharding + completeness marker, andUCVolumeWriter/UCVolumeReaderon AI Runtime.
Related¶
- systems/pytorch — the framework DCP ships in.
- distributed-checkpoint — the concept DCP implements.
- checkpoint-frequency — what cheap DCP saves unlock.
- concepts/durable-execution — the complementary data-side save.
- concepts/tensor-parallelism — the sharded regime DCP was built for.
- systems/unity-catalog-volumes — the AI Runtime storage target for DCP.
- asynchronous-checkpoint-staging — the
async_savestaging pattern. - resume-to-latest-complete-checkpoint — recovery via the
.metadatamarker.