The Three Pillars of Distributed GPU Training
Every distributed training strategy answers one question: how do you split the model across GPUs? The answer determines your memory footprint per GPU, communication volume between GPUs, and the utilization efficiency of your cluster. The three fundamental strategies are data parallelism, tensor parallelism, and pipeline parallelism.
Modern training at scale uses all three simultaneously in a 3D parallelism configuration. The art is in choosing the right ratios. A poorly balanced 3D configuration wastes GPU cycles on communication overhead or leaves memory underutilized. The optimal mix depends on the model architecture, the GPU interconnect topology, and the cluster size.
Data Parallelism: Simple Scaling with Limits
Data parallelism replicates the entire model on each GPU. Each GPU processes a different micro-batch of data and computes gradients independently. Gradients are averaged across all GPUs using an all-reduce operation, and the updated weights are broadcast back. This is the simplest strategy to implement and works well when the model fits entirely in a single GPU's VRAM.
The limitation is memory. A 70B parameter model in BF16 requires 140 GB just for parameters. The H200 has 141 GB of HBM3e. There is no room for optimizer states, activations, or the KV cache. Data parallelism alone cannot train models larger than approximately 30-40B parameters on current hardware. It also suffers from diminishing returns beyond 64 GPUs due to all-reduce synchronization overhead.
ZeRO (Zero Redundancy Optimizer) stages in DeepSpeed and FSDP extend data parallelism by sharding optimizer states, gradients, and parameters across GPUs while retaining the data-parallel programming model. Stage 3 ZeRO effectively becomes a form of parameter sharding that looks like data parallelism to the user but shares memory characteristics with model parallelism.
Tensor Parallelism: Splitting Within a Layer
Tensor parallelism (TP) splits individual weight matrices across GPUs. Instead of one GPU computing the full matrix multiplication for a transformer layer, the computation is distributed across multiple GPUs. The attention heads in multi-head attention are split across GPUs, and each GPU computes its subset of heads independently before an all-reduce combines the results.
TP requires high-bandwidth GPU interconnect because every transformer layer forward pass involves two all-reduce operations (one in the attention computation, one in the MLP). On NVLink-connected GPUs (900 GB/s on H200, 1.8 TB/s on B300), the communication overhead is manageable. On PCIe-only connections, TP becomes prohibitively slow for anything beyond 2 GPUs.
The advantage is that TP reduces per-GPU memory proportionally. A 70B model requires 140 GB of VRAM for parameters alone. With TP=8 across 8 GPUs, each GPU holds 17.5 GB of parameters. This enables fitting very large models that would never fit on a single GPU. The practical maximum for TP is typically 8 GPUs, matching the NVLink domain size of a single HGX baseboard.
Pipeline Parallelism: Layer-Based Sharding
Pipeline parallelism (PP) partitions the model by layer groups. GPU 0 handles layers 0-7, GPU 1 handles layers 8-15, and so on. Each GPU computes its layers and passes intermediate activations to the next GPU. This creates a sequential dependency: GPU 1 cannot start processing micro-batch N until GPU 0 finishes the forward pass and sends the activations.
The idle time caused by this sequential dependency is called the pipeline bubble. With a naive implementation, GPUs at the beginning and end of the pipeline spend half their time idle. The standard mitigation is to split each micro-batch into micro-micro-batches and overlap the forward and backward passes (the 1F1B schedule in PipeDream and Megatron-LM). With enough micro-batches, the bubble shrinks to approximately (p-1)/(m) where p is the pipeline depth and m is the number of micro-batches.
PP communicates only between adjacent pipeline stages, so it does not require full NVLink connectivity. A single PCIe link between nodes is sufficient for PP communication. This makes PP the strategy of choice for multi-node training where inter-node bandwidth is limited, provided the pipeline depth is kept small enough to avoid large bubbles.
Parallelization Strategy Comparison
The three strategies differ fundamentally in their communication pattern, memory profile, and scaling efficiency. Choosing the right one requires understanding where your specific bottleneck lies.
| Dimension | Data Parallelism | Tensor Parallelism | Pipeline Parallelism |
|---|---|---|---|
| Memory per GPU | Full model + opt | Model / TP size | Model / PP depth |
| Communication pattern | All-reduce (global) | All-reduce (within node) | P2P send/recv (adjacent) |
| Interconnect requirement | High (RDMA) | Very high (NVLink) | Moderate (PCIe is OK) |
| Scaling limit | ~64-128 GPUs | ~8 GPUs per node | ~16-32 pipeline stages |
| Pipeline bubble | None | None | Yes (bubble overhead) |
| Best for | Small models, large batch | Large models, fast interconnects | Multi-node, limited BW |
| Implementation | DDP, FSDP, ZeRO | Megatron-LM TP | Megatron-LM PP, GPipe |
Hybrid 3D Parallelism at Production Scale
A 1-trillion-parameter MoE model training on a 512-GPU H200 cluster might use: DP=8 across nodes, TP=8 within each node, PP=8 across nodes within each DP group. This gives 8 x 8 x 8 = 512 GPUs with balanced communication. The TP handles intra-node model sharding over NVLink, PP handles inter-node layer sharding over InfiniBand, and DP handles data replication with gradient synchronization over the same InfiniBand fabric.
The key tuning parameters are the TP degree (limited by the NVLink domain), the PP depth (limited by bubble overhead and micro-batch count), and the DP degree (limited by global batch size and convergence requirements). A 3D parallel training job on 512 GPUs might use a global batch size of 2,097,152 tokens with sequence parallelism enabled to fit activation memory.
Auto-tuning tools like DeepSpeed's autotune and Google's Pathways optimizer can search the configuration space automatically. For most teams, a sensible starting point is TP=8 (fill the NVLink domain), then choose the PP depth such that the pipeline bubble is under 10%, then use the remaining GPUs for data parallelism.
Memory and Communication Tradeoffs in Practice
For a concrete comparison: training Llama 3 70B on a 64-GPU cluster (8 nodes of 8x H100 SXM). With pure FSDP (data parallelism), each GPU holds the full model parameters (140 GB) plus optimizer states (280 GB with Adam) plus activations. This simply does not fit in 80 GB HBM. You are forced into model parallelism.
With TP=8 and PP=8 on the same 64 GPUs (removing DP), each GPU holds roughly 1/8 of the parameters and optimizer states but must communicate at every transformer layer through all-reduce (TP overhead) and P2P (PP overhead). The measured Model FLOPs Utilization (MFU) is typically 40-50%, compared to 50-60% for a well-tuned FSDP configuration on a smaller model that fits in memory.
The emerging consensus for 2026 training runs: use TP within the NVLink domain (8 GPUs), use PP across nodes within a rack (2-4 stages to keep bubble under 5%), and use DP for the remaining parallelism dimension. ZeRO-3 and FSDP2 have largely obsolete pure model parallelism for many workloads by providing a unified interface that dynamically switches between strategies per layer.
Choosing the Right Strategy for Your Workload
For models under 13B parameters that fit in a single GPU: use data parallelism (FSDP or DDP) and focus on scaling the data loading pipeline. For models 13B-70B that barely fit: use FSDP with ZeRO-3 or DeepSpeed stage 3 with CPU offloading for optimizer states. For models 70B-300B: use 3D parallelism with TP=8 within each node and PP=2-4 across nodes.
For models above 300B parameters or MoE architectures: use expert parallelism (a fourth strategy not covered here) alongside TP, PP, and DP. The Megatron-Core library from NVIDIA and DeepSpeed from Microsoft both provide production-tested implementations of all four strategies with automatic configuration search. The marginal cost of extra GPU hours spent tuning parallelism is recouped within days of production training at scale.
