Skip to content

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 .metadata records 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. .metadata is 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

  • 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_save staging pattern.
  • resume-to-latest-complete-checkpoint — recovery via the .metadata marker.
Last updated · 766 distilled / 2,225 read