Distributed Training Standard Architecture
Last reviewed: September 2026 | This is a fast-moving area subject to quarterly review.
Overview
Section titled “Overview”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.”
Data Parallelism (DP)
Section titled “Data Parallelism (DP)”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.
Tensor Parallelism (TP)
Section titled “Tensor Parallelism (TP)”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.
Pipeline Parallelism (PP)
Section titled “Pipeline Parallelism (PP)”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.
3D Hybrid Parallelism
Section titled “3D Hybrid Parallelism”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 |
Distributed Training Frameworks
Section titled “Distributed Training Frameworks”| 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 |
Checkpoint Strategy
Section titled “Checkpoint Strategy”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.
Related Documents
Section titled “Related Documents”- Next: Kubernetes, scheduling, gang scheduling — GPU Kubernetes and Scheduling
- Cluster communication tier, fabric, placement — GPU Workload Characteristics and Reference Architecture
- Automatic failure resume, capacity operations — Inference Serving, Reliability, and Cost
- Training pipeline within the AI system lifecycle — AI System Lifecycle and Engineering
Common Mistakes
Section titled “Common Mistakes”- 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
Checklist
Section titled “Checklist”- 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?