Skip to content

Distributed Training Standard Architecture

Last reviewed: September 2026 | This is a fast-moving area subject to quarterly review.

If the model and data fit on a single GPU, there’s nothing to worry about. The problem is that modern models far exceed a single GPU’s memory. So we split the work across multiple GPUs to train, and there are three main ways to split.

An analogy of several cooks sharing a large cooking job makes this easy to grasp.

  • Data parallelism (DP) — Each cook keeps a full copy of the recipe (the whole model) and only the customers (data) are divided; each cooks their share, then they reconcile the results.
  • Tensor parallelism (TP) — When one dish is too big for a single cook, several cooks make that one plate together at the same time.
  • Pipeline parallelism (PP) — Split the cooking stages so cook A preps, B grills, C plates — a relay-style hand-off.

Large-scale training layers all three together (3D parallelism), and each method trades off differently on “how often they must communicate / how much memory they save / how complex they are to implement.”

Replicate the whole model on each GPU, split only the data for processing, then synchronize the results (gradients) via all-reduce. It’s the simplest approach and the standard when the model fits on a single GPU. (Analogy: each cook holds a copy of the same recipe and only the customers are divided among them, then results are reconciled.)

  • Communication — Every training step, all GPUs combine their gradients. (all-reduce = a collective that gathers and sums every GPU’s value, then distributes the result back to all.)
  • Limitation — If the model itself exceeds a single GPU’s memory, this alone isn’t enough.
  • Saving memory (FSDP/ZeRO) — Shard the model’s parameters and intermediate state into small pieces spread across GPUs. This keeps data parallelism while training larger models beyond a single GPU’s limit.

Several GPUs shard the weights of a single layer (a computational layer that makes up the model) and compute it at the same time. Used when a single layer is too big to fit on one GPU.

  • Communication — GPUs exchange very frequently inside a layer, making it extremely latency-sensitive. So it is mostly used within one server (GPUs joined by NVLink).
  • Effect — Essential when a single layer exceeds GPU memory. Because communication is so frequent, extending it across nodes can cause throughput to collapse.

Divide the model’s layers into a few stages placed on different GPU groups, and stream the data as small pieces (micro-batches) through them like a relay.

  • Communication — Results are passed only at the boundary between stages, so communication volume is low. This makes it easy to spread across multiple servers.
  • Limitation — By its relay nature, idle gaps appear while waiting for the previous stage (pipeline bubbles); mitigate by increasing the number of data pieces.

Large-scale pre-training combines the three approaches hierarchically. Typically TP intra-node, PP inter-node, and DP on top.

Parallelism Communication volume Memory savings Latency sensitivity Recommended placement
Data parallel (DP) High (gradient all-reduce) None (large with FSDP/ZeRO) Medium Cluster-wide
Tensor parallel (TP) Very high (within a layer) Large Very high Intra-node (NVLink)
Pipeline parallel (PP) Low (stage boundaries) Large Low Inter-node
Framework Main parallelism Characteristics
PyTorch FSDP DP (sharding) PyTorch-native, distributes parameters/optimizer states
DeepSpeed DP(ZeRO) + PP + TP ZeRO staged memory optimization, offloading support
Megatron-LM TP + PP + DP TP implementation optimized for large-scale Transformer pre-training

Large-scale training runs for hours to weeks, so checkpointing to withstand node failures is essential. Checkpoint design is a matter of balancing save frequency against storage bandwidth.

  • Frequency — Too frequent, and save overhead leaves GPUs idle; too infrequent, and the compute lost on failure grows.
  • Storage bandwidth — Writing hundreds of GB to several TB of checkpoints in a short time makes storage-tier throughput a bottleneck.
  • Asynchronous/distributed saving — Save in the background without stopping training, or have each GPU save only its own shard in parallel to cut time.
  • Automatic resume — The flow of resuming from the last checkpoint after failure detection is covered in Inference Serving, Reliability, and Cost — Reliability and Failure Handling.
  • Introducing unnecessary 3D parallelism — Adding tensor/pipeline parallelism to a model well served by a single node, increasing only complexity and debugging cost
  • Extending tensor parallelism across nodes — Widening latency-sensitive TP beyond NVLink, creating a communication bottleneck
  • Only raising checkpoint frequency — Not scaling storage bandwidth accordingly, increasing GPU idle time during saves
  • Overlooking optimizer-state memory — Not accounting for the memory taken by optimizer states/gradients beyond parameters, causing OOM
  • Did you calculate whether the model and optimizer states fit in single-GPU/single-node memory?
  • If a single node suffices, did you first consider FSDP/ZeRO data parallelism?
  • Did you limit tensor parallelism to intra-node (NVLink) and place pipeline parallelism across nodes?
  • Did you design checkpoint frequency together with storage bandwidth?
  • Did you verify the automatic-resume flow on failure?