Distributed Data Parallelism (DDP) is the most widely-used distributed training strategy that replicates the entire model on every GPU and partitions the training data across GPUs — where each GPU computes gradients on its data partition and then all GPUs synchronize gradients via all-reduce before applying the same parameter update, ensuring all replicas remain identical while achieving near-linear throughput scaling with the number of GPUs.
How DDP Works
1. Initialization: The model is replicated identically on N GPUs. Each GPU receives a different shard of the training data (via DistributedSampler). 2. Forward Pass: Each GPU computes the forward pass on its local mini-batch independently. 3. Backward Pass: Each GPU computes gradients on its local mini-batch. Gradients are different on each GPU (different data). 4. All-Reduce: Gradients are summed (and averaged) across all GPUs using an efficient collective operation (NCCL ring or tree all-reduce). After all-reduce, every GPU has identical averaged gradients. 5. Parameter Update: Each GPU applies the identical optimizer step using the identical averaged gradients, maintaining weight synchrony.
Scaling Behavior
- Throughput: Near-linear scaling — N GPUs process N mini-batches per step. Effective batch size = per-GPU batch × N.
- Communication Overhead: All-reduce transfers 2 × model_size bytes per step (for a ring all-reduce). For a 7B parameter model in FP16/BF16: 2 × 14 GB = 28 GB of all-reduce traffic per step.
- Computation-Communication Overlap: PyTorch DDP and DeepSpeed overlap the all-reduce of early layers' gradients with the backward pass of later layers. This hides most of the communication latency behind useful compute.
Large Batch Training Challenges
- Learning Rate Scaling: Linear scaling rule — multiply the base learning rate by N (GPUs). Works up to a point; very large batch sizes (>32K) require warm-up and special optimizers (LARS, LAMB).
- Generalization Gap: Extremely large batch sizes can degrade model quality (sharper minima). Gradient noise reduction at large batch sizes reduces the implicit regularization of SGD.
- Batch Normalization: BN statistics computed per-GPU with small local batch sizes are noisy. SyncBatchNorm computes statistics across all GPUs but adds communication overhead.
Implementations
- PyTorch DDP:
torch.nn.parallel.DistributedDataParallel. Wraps any model, handles gradient synchronization transparently via NCCL backend. Supports gradient accumulation for effective batch size scaling without more GPUs. - DeepSpeed ZeRO: Extends DDP by partitioning optimizer states (ZeRO-1), gradients (ZeRO-2), and parameters (ZeRO-3) across GPUs, reducing per-GPU memory. Enables training models that don't fit in a single GPU's memory while maintaining data-parallel semantics.
- Horovod: Framework-agnostic distributed training library.
hvd.DistributedOptimizerwraps any optimizer with all-reduce gradient synchronization.
Distributed Data Parallelism is the workhorse of large-scale model training — the strategy that scaled deep learning from single-GPU research experiments to thousand-GPU production training runs by distributing the data while keeping the model replicated and synchronized.
Explore 500+ Semiconductor & AI Topics
From EUV lithography to CUDA optimization — search the full knowledge base or chat with our AI assistant.