THE PARALLELISM TAXONOMY: 5 STRATEGIES IN ONE FRAMEWORK
Modern distributed training combines up to five parallelism strategies. Data parallelism (DP) replicates the entire model on each GPU and processes different micro-batches; FSDP is the dominant DP variant that shards parameters across replicas. Tensor parallelism (TP) splits individual layers' weight matrices across GPUs, requiring an all-reduce after each layer. Pipeline parallelism (PP) partitions model layers across GPUs, with each GPU computing a subset of layers sequentially. Sequence parallelism (SP) splits the sequence dimension across GPUs within attention computation. Expert parallelism (EP) distributes MoE expert modules across GPUs.
Each strategy targets a different bottleneck. DP maximizes compute utilization but has communication overhead from gradient all-reduce. TP reduces per-GPU memory by sharing parameters but adds communication on every layer's forward and backward pass (all-reduce of activations). PP reduces communication volume to only the layer boundaries (point-to-point sends) but introduces idle bubbles. SP addresses the O(n^2) memory of attention at long sequences. EP addresses the expert capacity problem in MoE models. The optimal combination depends on model size, sequence length, GPU count, and inter-GPU bandwidth.
| Strategy | Memory Redux | Communication | GPU Idle Time | Scaling Limit | Best For |
|---|---|---|---|---|---|
| DP / FSDP | 3-10x (ZeRO-3) | Grad all-reduce / param gather | None (sync) | 1024+ GPUs | Small models, high throughput |
| TP (tensor) | Nx layer weight | Attn/MLP all-reduce each layer | None | 8-16 GPUs | Large models, single node |
| PP (pipeline) | Nx total weights | Layer boundary p2p only | 10-30% bubble | 16-128 GPUs | Cross-node, low bandwidth |
| SP (sequence) | Nx attn memory | Attn all-reduce per layer | None | 8-64 GPUs | Long sequences (128K+) |
| EP (expert) | Nx expert weights | Token dispatch all-to-all | Load imbalance | 64-1024 GPUs | MoE models |
FSDP: THE DEFAULT DATA PARALLELISM STRATEGY
FSDP (Fully Sharded Data Parallelism) shards model parameters, gradients, and optimizer states across data-parallel ranks. At each forward pass, FSDP all-gathers the parameters for the current layer, computes, then discards non-current-layer parameters. The backward pass repeats the all-gather and adds a reduce-scatter for gradients. Communication volume per layer is 2x the parameter size (one all-gather for forward, one all-gather + one reduce-scatter for backward). FSDP-3 (full sharding) achieves near-perfect scaling up to 64 GPUs for moderate-sized models (7B-13B).
Beyond 64 GPUs, FSDP scaling efficiency degrades to 70-85 percent due to the increasing all-gather latency. The all-gather is a blocking collective operation: all GPUs must complete it before any GPU can continue computation. At 256 GPUs, the all-gather for a single transformer layer's parameters (14M for Llama-70B attention) takes 50-100 microseconds on NVLink but 200-500 microseconds on InfiniBand across nodes. With 80 layers, this adds 16-40 ms of communication overhead per step, reducing throughput by 15-25 percent versus ideal linear scaling.
The FSDP configuration hierarchy: FSDP-1 (optimizer state sharding) for up to 128 GPUs, FSDP-2 (gradient sharding) for up to 256 GPUs, FSDP-3 (full sharding) for up to 512 GPUs. Beyond 512 GPUs, FSDP-3 with hybrid sharding (shard within nodes, replicate across nodes) or TP+PP+DP combination is necessary.
| Model | FSDP-1 (GPUs) | FSDP-2 (GPUs) | FSDP-3 (GPUs) | FSDP-3 Efficiency | Recommended |
|---|---|---|---|---|---|
| 1B | 16 | 32 | 64-128 | 90-95% | FSDP-1 or DDP |
| 7B | 16-32 | 32-64 | 128-256 | 85-92% | FSDP-2 |
| 13B | 16-32 | 64-128 | 128-256 | 82-90% | FSDP-2 + TP=2 |
| 70B | 32-64 | 64-128 | 256-512 | 75-85% | FSDP-3 + TP=4 |
| 405B | 64-128 | 128-256 | 512-1024 | 65-78% | TP=8 + PP=4 + DP |
TENSOR AND PIPELINE PARALLELISM: THE COMMUNICATION TRADEOFF
Tensor parallelism (TP) splits each weight matrix column-wise or row-wise across GPUs. For a 70B model with TP=8, each GPU holds 1/8 of every layer's weights. Every transformer layer requires 6 all-reduce operations per forward pass (one for QKV projection, one for attention output, two for MLP gate+up, one for MLP down, one for final output). At 70B with 80 layers, each all-reduce communicates 1/8 of the layer's activation tensor - for a 4K sequence length with hidden dimension 8192, this is 4K x 1K = 4M elements per all-reduce, or 8 MB at FP16. Six all-reduces per layer x 80 layers = 480 all-reduces per step, totaling 3.8 GB of cross-GPU traffic.
Pipeline parallelism (PP) divides layers across GPUs. With PP=8 on a 80-layer model, each GPU holds 10 layers. Communication is one point-to-point send per layer boundary per micro-batch: 8 MB per activation tensor (same as the all-reduce size in TP) per forward and backward pass. Total communication per step: 2 x 8 MB x (PP-1) x num_microbatches. With 32 micro-batches and PP=8, this is 2 x 8 x 7 x 32 = 3.6 GB - similar to the TP volume. The critical difference is that PP communication is asynchronous (send while computing the next micro-batch), hiding latency behind computation, while TP requires synchronous all-reduce at every layer.
TP achieves zero bubble (no idle GPU time), while PP introduces a pipeline bubble that wastes 10-30 percent of compute. The 1F1B (one-forward-one-backward) scheduling reduces the bubble to (PP-1) / (num_microbatches) fraction. With PP=8 and 32 micro-batches, bubble overhead is 7/32 = 22 percent. Interleaved 1F1B (multiple layers per stage) reduces this to 7/64 = 11 percent.
| Configuration | Comm per Step | Bubble Overhead | Memory per GPU | 70B Throughput | Best Network |
|---|---|---|---|---|---|
| FSDP-3 only | ~2.5 GB | 0% | 18 GB | 1.0x (baseline) | NVLink |
| TP=8, DP=8 | ~3.8 GB | 0% | 10 GB | 0.85-0.92x | NVLink required |
| PP=8, DP=32 | ~3.6 GB | 22% | 18 GB | 0.75-0.80x | Any interconnect |
| TP=4, PP=4, DP=16 | ~2.2 GB TP + ~1.0 GB PP | 11% | 16 GB | 0.88-0.92x | NVLink in node |
| TP=8, PP=4, DP=8 | ~3.0 GB TP + ~0.8 GB PP | 11% | 9 GB | 0.82-0.88x | NVLink in node |
| SP=8 (long ctx) | ~2 GB per layer | 0% | Seq_len/8 attn | 0.75-0.85x | NVLink |
EXPERT PARALLELISM: MIXTURE-OF-EXPERTS DISTRIBUTION
Expert parallelism (EP) distributes MoE expert modules across GPUs, with a router dispatching each token to its designated expert GPU. The key infrastructure challenge is the all-to-all communication: each GPU sends tokens to the GPU hosting each expert and receives back the expert's output. For a Mixtral 8x7B MoE with 8 experts across 8 GPUs, each step requires an all-to-all of approximately batch_size x hidden_size x 2 (send + receive) per layer. At batch size 32 with 4K sequence length and hidden size 4096, each all-to-all communicates 32 x 4K x 4K x 2 bytes = 1 GB per MoE layer. With 32 MoE layers, that's 32 GB per step - an order of magnitude more communication than dense TP.
Expert load imbalance creates GPU idle time. The router distributes tokens to experts based on the gating function, but the distribution is rarely uniform. A 10-20 percent imbalance means some GPUs process fewer tokens and sit idle waiting for the straggler. The standard mitigation is auxiliary load-balancing loss (2-5 percent weight) that penalizes the router for imbalanced assignments. With proper tuning, imbalance drops to 2-5 percent, but the loss gradient adds 5-10 percent to training FLOPs. DeepSpeed-MoE implements a dynamic expert placement strategy that re-assigns experts to GPUs during training to minimize cross-node communication for expert pairs that frequently co-activate.
THE COMBINED STRATEGY: 3D PARALLELISM IN PRACTICE
Production training at 70B+ scale uses 3D parallelism: TP within a node (NVLink required), PP across nodes (InfiniBand), and DP across the cluster. The canonical configuration for 70B on 64 H100 GPUs (8 nodes x 8 GPUs): TP=8 (within node), PP=2 (2 stages across 4 nodes each), DP=4. This configuration achieves 220-250 teraFLOP/s per GPU (50-55 percent utilization) on H100 for the 70B model with 4K sequence length.
The scaling to 405B models on 512 GPUs: TP=8 (within node), PP=8 (node-level stages), DP=8. At 8K sequence length, this achieves 180-210 teraFLOP/s per GPU. The 28 percent lower utilization versus 70B is due to increased PP bubble (PP=8 has higher bubble than PP=2) and higher communication-to-computation ratio for the larger model. The 405B configuration requires 4.1 TB total HBM across 512 H100 GPUs - just barely fitting the 2.7 TB of weights + optimizer states + activations.
The trend for 2025-2026 is toward sequence parallelism for long-context training and expert parallelism for MoE models. Megatron-LM and NVIDIA NeMo both support SP=8 for 1M+ context training, while DeepSpeed supports EP mixed with TP/PP. The practical recommendation: start with FSDP-3, add TP when FSDP communication efficiency drops below 80 percent (typically at 70B+ on 256+ GPUs), add PP when TP needs exceed 8 GPUs, and add SP when sequence length exceeds 128K.
| Model | GPUs | TP | PP | DP | SP | TFLOPS/s per GPU | Utilization |
|---|---|---|---|---|---|---|---|
| 7B | 16 | 1 | 1 | 16 (FSDP-2) | 1 | 380-420 | 85-93% |
| 13B | 32 | 2 | 1 | 16 | 1 | 350-390 | 78-87% |
| 70B | 64 | 8 | 2 | 4 | 1 | 220-250 | 50-55% |
| 70B (128K ctx) | 128 | 4 | 4 | 4 | 4 | 180-210 | 40-47% |
| 405B | 512 | 8 | 8 | 8 | 1 | 180-210 | 40-47% |
| 405B (MoE) | 256 | 4 | 4 | 4 | 1 | 250-290 | 56-64% |
B200: SIMPLIFYING THE PARALLELISM STACK
B200's 2.4x memory capacity and 1.8x memory bandwidth over H100 fundamentally change parallelism strategy. A 70B model fits on 4 B200 GPUs with FSDP-3 and no TP/PP, versus 8 H100 GPUs requiring TP=8. This reduces TP's synchronous all-reduce communication by 50 percent, increasing per-GPU utilization to 65-75 percent. For 405B MoE models (e.g., DeepSeek-V3 style with 37B active params), B200 enables single-node (8 GPU) deployment that required 2-4 nodes on H100.
The practical impact is that models up to 70B require only FSDP-3 on B200, eliminating the complexity of 3D parallelism configuration that requires significant tuning expertise. For 405B dense models, the B200 reduces the required GPU count from 512 to 256-320, proportionally reducing the distributed communication overhead. The simplified parallelism stack means faster training startup, fewer communication hangs, and easier debugging - benefits that compound over multi-month training runs.
