Data Parallelism: DDP, Gradient Synchronization & All-Reduce

Michael BrenndoerferJanuary 16, 202661 min read

Part of Language AI Handbook

Explains how DDP trains large models across multiple GPUs using ring all-reduce gradient synchronization, gradient bucketing.

Choose your expertise level to adjust how many terms are explained. Beginners see more tooltips, experts see fewer to maintain reading flow. Hover over underlined terms for instant definitions.

Article links

Make inline references clickable

Data Parallelism

Training a large language model on a single GPU is impractical. The datasets are enormous, the models are massive, and a single forward-backward pass can take minutes. The first and most natural response to this problem is to ask: what if we simply used more GPUs running the same computation at the same time?

This is the core idea behind data parallelism. Instead of changing the model or splitting it across devices, you keep a complete copy of the model on every GPU and divide the training data among them. Each GPU independently computes gradients on its assigned batch, and then all GPUs exchange those gradients so that every copy of the model stays synchronized. The result is that you can process NN times as much data per second with NN GPUs, cutting wall-clock training time roughly by a factor of NN.

Data parallelism is the oldest and most widely used form of distributed training. Before ZeRO, tensor parallelism, and pipeline parallelism became mainstream (topics we will cover in the next several chapters), data parallelism was the default strategy for scaling language model training. It remains the backbone of distributed training even when combined with other parallelism strategies. Systems like GPT-3, PaLM, LLaMA, and virtually every production language model training run in recent years have relied on data parallelism as their primary scaling axis.

The appeal of data parallelism is its simplicity. Conceptually, you are just running the same training code on multiple GPUs with different data. You do not need to redesign your model architecture, change your optimizer, or reason about complicated data flow between different parts of the model. The distribution is transparent: from the model's perspective, it is receiving gradients and applying updates just like single-GPU training. The parallelism lives entirely in the training harness.

This chapter covers the full lifecycle of data parallelism: the DDP algorithm that makes it work, the gradient synchronization protocol that keeps model replicas in sync, the all-reduce operation that powers that synchronization, the scaling behaviour you can realistically expect as you add more GPUs, and the practical configuration decisions you need to make when setting up a real distributed training job. By the end, you will understand how to use PyTorch DDP and why each of its design decisions exists.

The Data Parallelism Problem

Before we look at how data parallelism works, it helps to understand what problem it is solving and why naive approaches fail.

Suppose you have a model with parameters θ\theta and you want to train on a dataset D\mathcal{D} with NN GPUs. The gradient on the full dataset is the average over all examples:

∇L(θ)=1∣D∣∑x∈D∇Lx(θ)\nabla \mathcal{L}(\theta) = \frac{1}{|\mathcal{D}|} \sum_{x \in \mathcal{D}} \nabla \mathcal{L}_x(\theta)

where:

  • θ\theta is the vector of all model parameters
  • D\mathcal{D} is the full training dataset
  • Lx(θ)\mathcal{L}_x(\theta) is the loss on a single example xx
  • ∇Lx(θ)\nabla \mathcal{L}_x(\theta) is the per-example gradient, a vector with as many dimensions as there are parameters

If you split D\mathcal{D} into NN equal shards and compute the gradient on each shard separately, the average of those shard gradients equals the full-dataset gradient exactly. This is the mathematical foundation of data parallelism: gradient averaging over shards is equivalent to gradient computation over the union of those shards.

The challenge is how to perform that averaging efficiently across NN GPUs that are connected by a communication network.

Why the Parameter Server Fails

The simplest parallelism strategy would be to split D\mathcal{D} into NN shards, run each on a separate GPU, collect all gradients on one designated "parameter server" GPU, average them there, and broadcast the updated parameters back to all GPUs. This approach, called parameter server training, was popular in early distributed machine learning research (Dean et al., 2012) and is straightforward to implement.

Parameter servers have a fatal flaw that becomes apparent at scale: the parameter server is a communication bottleneck. Consider the data flow. Each of the NN worker GPUs must send its full gradient vector to the parameter server. The parameter server must receive NN gradient vectors, sum them, and then broadcast the result to all NN workers. The parameter server's inbound bandwidth must handle NN gradient vectors simultaneously, and its outbound bandwidth must serve NN simultaneous downloads.

For a 7B parameter model with BF16 gradients (2 bytes per parameter), the gradient tensor is 14 GB. With 64 GPUs, the parameter server must receive 64×14=89664 \times 14 = 896 GB of gradient data and send back 14 GB to each of 64 workers (another 896 GB) every single training step. No single GPU has that kind of network bandwidth. Even on InfiniBand HDR with 200 Gb/s links, transferring 896 GB would take over 35 seconds per step, compared to a compute time of perhaps 250 ms. The ratio is roughly 140-to-1 in favor of the bottleneck.

Asynchronous parameter servers alleviate the synchronization requirement but introduce gradient staleness: workers apply gradients computed on stale parameters, which can harm convergence. The tradeoffs are poor.

Modern data parallelism solves the bottleneck problem with a different communication pattern called all-reduce, which distributes the aggregation work across all GPUs equally so that no single GPU is a bottleneck.

The Correctness Requirement

One subtle point before we proceed: for data parallelism to produce the correct result, every GPU must apply the same parameter update at the end of each step. If different GPUs apply different updates (because they each used only their local gradients), the model replicas would diverge. After a few steps, you would have NN different models rather than one model trained on NN times as much data.

This is why gradient synchronization is so important. Before any GPU applies an optimizer update, all GPUs must have agreed on a single averaged gradient. The moment any GPU deviates from this shared gradient, the replicas diverge and training becomes incoherent. DDP's design makes this guarantee by construction: no optimizer step can proceed until the all-reduce is complete on every rank.

Distributed Data Parallel: The Algorithm

Distributed Data Parallel (DDP), introduced by Li et al. (2020) and implemented in PyTorch's torch.nn.parallel.DistributedDataParallel, is the standard approach for data-parallel training today. DDP was a significant improvement over PyTorch's earlier DataParallel module, which used a parameter-server-like approach on a single node and suffered from exactly the bottleneck problems described above. DDP's key properties are:

  • Each GPU (called a "rank") holds a complete copy of the model
  • Each rank processes a different shard of the mini-batch
  • Gradients are synchronized via all-reduce before the optimizer step
  • All ranks apply identical parameter updates, so all copies stay identical

Let us walk through one training step in detail, because each phase has interesting properties worth understanding.

Step 1: Forward Pass

Each rank rr receives a different mini-batch Br\mathcal{B}_r of size bb. The global effective batch size is N×bN \times b, where NN is the number of ranks. Each rank computes its local loss:

Lr=1∣Br∣∑x∈BrLx(θ)\mathcal{L}_r = \frac{1}{|\mathcal{B}_r|} \sum_{x \in \mathcal{B}_r} \mathcal{L}_x(\theta)

where:

  • Lr\mathcal{L}_r is the local loss on rank rr, a scalar
  • Br\mathcal{B}_r is the mini-batch assigned to rank rr, a set of examples
  • θ\theta is the current parameter vector, identical on all ranks at the start of the step

This forward pass is entirely local. No inter-GPU communication happens during the forward pass. Each rank can run its forward computation at full speed without waiting for any other rank. This is one of the reasons data parallelism is so efficient: the most compute-intensive phase of training (the forward pass through a large transformer) is completely parallel.

The forward passes on different ranks are not strictly synchronized either. If rank 2 finishes its forward pass slightly before rank 1, it can immediately start backpropagation without waiting. Synchronization only becomes necessary when gradients are ready and need to be averaged.

Step 2: Backward Pass with Gradient Bucketing

During backpropagation, each rank computes local gradients ∇rLr\nabla_r \mathcal{L}_r with respect to every parameter. Backpropagation visits layers in reverse order, starting from the output layer and working back toward the input embeddings. As each layer's backward pass completes, its parameter gradients become available.

DDP does not wait for the entire backward pass to finish before communicating. Instead, it overlaps gradient communication with backward computation using a technique called gradient bucketing. Understanding why this matters requires thinking about the timeline of a training step.

In a naive implementation, the full sequence of events would be: run the entire forward pass, run the entire backward pass (computing all gradients), communicate all gradients via all-reduce, then run the optimizer step. The communication phase sits between backpropagation and the optimizer update and cannot begin until all gradients are ready. During the communication phase, the GPU's compute units sit idle.

With gradient bucketing, this idle time is eliminated. Here is how it works. DDP groups parameters into buckets, with each bucket containing roughly 25 MB of gradients by default. The buckets are ordered in reverse layer order, so the last layers (whose gradients become available first during backward) form the first buckets. As soon as all gradients in a bucket are ready, DDP immediately launches an all-reduce operation on that bucket, while the backward pass continues computing gradients for earlier layers in parallel.

By the time backpropagation reaches the earliest layers (embeddings, first transformer blocks), several all-reduce operations have already completed for the later layers. When the backward pass finally finishes, the last all-reduce is either already done or nearly done. The total step time is reduced from compute + communication to approximately max(compute, communication) plus a small serial residual.

The choice of 25 MB per bucket involves a tradeoff. Smaller buckets reduce the latency before the first all-reduce starts (because buckets fill up faster), but they increase the total number of collective communication operations, each with fixed launch overhead. Larger buckets amortize that overhead but delay the start of communication and require more memory in the communication buffers. For most language model training, 25 MB is a reasonable default, but you can tune it with the bucket_cap_mb parameter.

Out[3]:
Visualization
Gantt chart comparing naive DDP and bucketed DDP training step timelines.
Training step timeline comparing naive DDP (top) versus DDP with gradient bucketing (bottom) for a model with 3 gradient buckets. In the naive approach, all communication happens after the backward pass completes, serializing compute and communication. With gradient bucketing, each bucket's all-reduce launches as soon as its gradients are ready, overlapping communication for later buckets with continued backward computation for earlier layers. The bucketing approach reduces total step time significantly when communication time is comparable to compute time.

Step 3: All-Reduce Gradient Synchronization

The all-reduce operation takes a tensor from every rank, computes an element-wise sum (or average), and writes the result back to every rank simultaneously. After all-reduce on gradients:

∇sync=1N∑r=0N−1∇rLr\nabla_{\text{sync}} = \frac{1}{N} \sum_{r=0}^{N-1} \nabla_r \mathcal{L}_r

where:

  • ∇sync\nabla_{\text{sync}} is the synchronized gradient, identical on all ranks after the operation completes
  • NN is the number of ranks (GPUs)
  • ∇rLr\nabla_r \mathcal{L}_r is the local gradient computed by rank rr, a vector of the same shape as θ\theta

This synchronized gradient is mathematically equivalent to the gradient you would have computed if you had processed all NN mini-batches on a single device. The averaging ensures that the gradient scale does not grow with the number of GPUs, so the same learning rate works regardless of NN (assuming the global batch size is proportionally increased). This is a critical design choice: without the averaging, gradients would grow proportionally to NN, requiring you to scale down the learning rate by NN every time you added more GPUs, which would break any learning rate schedule you had tuned.

The all-reduce operation is a collective operation, meaning it requires all ranks to participate simultaneously. Every rank must call all-reduce at the same time (with the same tensor shape), and every rank blocks until the operation completes. This synchronization barrier is what guarantees that all ranks end up with identical gradients.

Step 4: Optimizer Step

After gradient synchronization, every rank has identical ∇sync\nabla_{\text{sync}}. Each rank independently runs its optimizer (Adam, AdamW, SGD, or any other) using this synchronized gradient:

θ←θ−α⋅Optimizer(∇sync)\theta \leftarrow \theta - \alpha \cdot \text{Optimizer}(\nabla_{\text{sync}})

where α\alpha is the learning rate. Because every rank starts the step with identical θ\theta and identical ∇sync\nabla_{\text{sync}}, every rank computes an identical update and every rank ends the step with identical θ\theta. No parameter synchronization is needed after the optimizer step. The model copies remain perfectly in sync as long as they were in sync at the start, which is guaranteed by the all-reduce.

This design is elegant. You never need to communicate parameters. You only need to communicate gradients (once per step), and the rest of the training loop is identical to single-GPU training. The optimizer, gradient clipping, learning rate schedule, and all other components of the training pipeline require no modification.

Why Not Synchronize Parameters Instead of Gradients?

You might wonder why DDP synchronizes gradients rather than just synchronizing parameters after the optimizer step. Gradient synchronization gives mathematically exact results: all ranks compute the same gradient, apply the same update, and arrive at the same parameters. Parameter synchronization would also work correctly, but it would require synchronizing after the optimizer step rather than before it. This matters for gradient bucketing: DDP can overlap gradient communication with backpropagation (launching all-reduce as gradients become ready), but you cannot overlap parameter communication with anything because parameters are only ready after the optimizer finishes. ZeRO (covered in a later chapter) does synchronize parameters in some configurations, but as part of a more complex protocol that achieves memory savings unavailable with pure gradient synchronization.

The All-Reduce Operation

All-reduce is the communication primitive at the heart of DDP. Understanding how it works explains both why DDP is efficient and what its scaling limits are. All-reduce is not a single algorithm but a family of algorithms with different performance characteristics. Choosing the right one for your hardware topology can make a substantial difference in training speed.

The Reduce-Broadcast Baseline

The simplest correct all-reduce is a two-step reduce-broadcast:

  1. One designated root rank aggregates all tensors from all other ranks (reduce to root)
  2. The root broadcasts the result to all other ranks

This approach has the same parameter server bottleneck problem we described earlier. The root rank must receive data from N−1N-1 other ranks (inbound traffic =(N−1)⋅∣g∣= (N-1) \cdot |\mathbf{g}| bytes) and send data to N−1N-1 ranks (outbound traffic =(N−1)⋅∣g∣= (N-1) \cdot |\mathbf{g}| bytes). The root's required bandwidth scales linearly with NN. For large NN, this approach is impractical.

Ring All-Reduce

The standard solution is ring all-reduce, popularized by Baidu Research (Gibiansky, 2017) for deep learning. In ring all-reduce, the NN ranks are arranged in a logical ring topology. The algorithm has two phases, each consisting of N−1N-1 communication steps.

Phase 1: Reduce-Scatter

In the reduce-scatter phase, each rank splits its gradient tensor gr\mathbf{g}_r into NN equal chunks: gr=[gr(0),gr(1),…,gr(N−1)]\mathbf{g}_r = [\mathbf{g}_r^{(0)}, \mathbf{g}_r^{(1)}, \ldots, \mathbf{g}_r^{(N-1)}]. Over N−1N-1 communication steps:

  • Each rank sends one chunk to its right neighbor in the ring
  • Each rank receives one chunk from its left neighbor and adds it to its local copy of that chunk

After N−1N-1 steps, each rank holds one chunk that contains the sum of that chunk across all ranks. Specifically, rank rr holds ∑i=0N−1gi(r)\sum_{i=0}^{N-1} \mathbf{g}_i^{(r)}, the fully summed version of chunk rr.

Phase 2: AllGather

In the all-gather phase, each rank sends its fully-summed chunk to its right neighbor over N−1N-1 steps. By the end:

  • Every rank holds every chunk of the fully summed gradient
  • The complete synchronized gradient ∑rgr\sum_{r} \mathbf{g}_r is on every rank

The total data transferred by each rank across both phases is:

Data per rank=2⋅N−1N⋅∣g∣\text{Data per rank} = 2 \cdot \frac{N-1}{N} \cdot |\mathbf{g}|

where:

  • ∣g∣|\mathbf{g}| is the total size of the gradient tensor in bytes
  • N−1N\frac{N-1}{N} approaches 1 as NN grows, so each rank transfers approximately 2∣g∣2|\mathbf{g}| bytes regardless of the number of GPUs
  • The factor of 2 accounts for both phases: one reduce-scatter send and one all-gather send, each of approximately ∣g∣|\mathbf{g}| bytes total

The defining property of ring all-reduce is that the per-rank communication volume is independent of NN. Adding more GPUs does not increase the amount of data each GPU must send or receive. This is what makes ring all-reduce bandwidth-optimal: it achieves the theoretical minimum data transfer for the all-reduce problem.

To understand why this is remarkable, compare it to the naive reduce-broadcast: a single root rank must send and receive approximately N⋅∣g∣N \cdot |\mathbf{g}| bytes, which grows linearly with NN. Ring all-reduce eliminates this growth by distributing the aggregation work evenly.

Out[4]:
Visualization
Diagram of ring all-reduce reduce-scatter phase with 4 GPUs and clockwise arrows.
Ring all-reduce reduce-scatter phase for 4 GPUs. After N-1 steps, each GPU holds the fully summed gradient for its assigned chunk, with GPU 0 owning sum[0], GPU 1 owning sum[1], and so on. The clockwise arrows show the flow of one chunk per step.
Diagram of ring all-reduce all-gather phase with 4 GPUs holding the full gradient.
Ring all-reduce all-gather phase. Starting from the partial sums computed in reduce-scatter, each GPU propagates its summed chunk clockwise until every GPU holds the complete gradient. Each GPU transfers approximately 2 x gradient bytes total across both phases, independent of GPU count.

Tree All-Reduce

Ring all-reduce is optimal for large messages where bandwidth is the bottleneck. For small messages (small models, small gradient tensors), the latency of 2(N−1)2(N-1) communication steps becomes significant even if each step transfers little data. In this regime, a tree topology performs better.

In a binary tree all-reduce, ranks are organized as a binary tree. The reduce phase propagates data up the tree from leaves to root, with each internal node summing the values it receives from its children. The broadcast phase propagates the result back down from root to leaves. The total number of communication steps is 2log⁡2N2 \log_2 N rather than 2(N−1)2(N-1) steps, which is dramatically fewer for large NN.

The tradeoff is bandwidth: the root of the tree must process and retransmit O(log⁡N)O(\log N) times as much data as leaf nodes, creating mild hotspots. For the large gradient tensors typical in language model training, ring all-reduce generally wins on throughput because bandwidth, not latency, is the bottleneck. For very small tensors (such as bias gradients in individual layers), tree-based all-reduce can be faster.

NCCL and Hardware Topology

In practice, PyTorch uses NVIDIA Collective Communications Library (NCCL) rather than implementing ring all-reduce directly in Python or CUDA. NCCL is topology-aware: it detects whether GPUs are on the same PCIe bus, connected via NVLink, or spread across multiple nodes connected by InfiniBand or RoCE, and chooses the optimal communication algorithm for each case. NCCL also handles buffer management, stream synchronization, and error recovery, all of which are tedious to implement correctly.

NVLink (the high-bandwidth GPU interconnect available on NVIDIA A100 and H100 servers) provides roughly 600 GB/s aggregate bidirectional bandwidth on modern 8-GPU nodes, compared to roughly 32 GB/s for PCIe. For multi-node training, InfiniBand HDR provides roughly 200 Gb/s per link (about 25 GB/s). The bandwidth difference between NVLink and InfiniBand is enormous: NVLink is roughly 24 times faster per GPU. This means that gradient communication that takes 20 ms on NVLink takes nearly 500 ms on InfiniBand, completely dominating compute time.

The topology distinction is particularly important for multi-node training. Within a single node with 8 GPUs connected via NVLink, NCCL uses a ring or tree topology that fully exploits the high NVLink bandwidth. Across nodes, NCCL routes traffic through the node's InfiniBand host channel adapter (HCA), which has far lower bandwidth per GPU than NVLink. A typical multi-node DDP job therefore has a two-level communication hierarchy: intranode communication using the fast NVLink fabric, and internode communication using the slower InfiniBand fabric.

Understanding this hierarchy helps you make informed decisions about how to partition GPUs. If you have a cluster of 8-GPU nodes, it is almost always better to fill entire nodes before spreading to additional nodes. Keeping communication within a single node maximizes the use of NVLink and minimizes InfiniBand traffic. Once you must cross nodes, you want to keep the cross-node gradient transfers as small as possible, which is one motivation for combining DDP with ZeRO (which reduces the gradient tensor that must be communicated across nodes).

NCCL vs. Gloo Backends

PyTorch's dist.init_process_group supports multiple backends. NCCL is the correct choice for GPU training: it is tightly integrated with NVIDIA hardware and provides the best performance. Gloo is the alternative backend that runs on CPU and works across platforms including macOS and Windows. Gloo is useful for debugging and for single-machine multi-process experiments where GPU communication is unavailable. MPI is a third option for clusters that have MPI already configured, but NCCL is generally preferred for pure GPU training. When you see code examples using "gloo", that is always for portability or demonstration, not for production performance.

DDP in Practice: PyTorch Implementation

Let us walk through a complete DDP training setup using PyTorch. Modern PyTorch DDP requires specific components to work correctly, and the setup is slightly more involved than single-GPU training.

Launching with torchrun

Before looking at the code, it is worth understanding how DDP processes are launched. DDP requires one process per GPU, and each process must know its rank (which GPU it is), the total number of processes (world size), and how to find the other processes (via a rendezvous point, typically a hostname and port).

PyTorch provides torchrun (replacing the older torch.distributed.launch) as the recommended launcher. You invoke it as:

In[8]:
Code
# torchrun --nproc_per_node=8 --nnodes=2 --node_rank=0 \
#     --master_addr=node0 --master_port=29500 train.py

torchrun starts the specified number of processes on each node, sets the environment variables RANK, LOCAL_RANK, WORLD_SIZE, MASTER_ADDR, and MASTER_PORT, and waits for all processes to join the rendezvous before any of them proceed. Your training script reads these environment variables to configure its local rank and total world size.

For single-node training, the simplest invocation is:

In[10]:
Code
# torchrun --nproc_per_node=8 train.py

This starts 8 processes, one per GPU, with ranks 0 through 7.

Setting Up the Process Group

In[5]:
Code
import os

import torch.distributed as dist


def setup_distributed():
    """Initialize the distributed process group."""
    # In real training, these are set by the launcher (torchrun or SLURM)
    # We simulate a single-process environment here
    if not dist.is_initialized():
        os.environ.setdefault("MASTER_ADDR", "localhost")
        os.environ.setdefault("MASTER_PORT", "12355")
        dist.init_process_group(
            backend="gloo",  # Use "nccl" for real GPU training
            rank=0,
            world_size=1,
        )


def cleanup_distributed():
    """Clean up the distributed process group."""
    if dist.is_initialized():
        dist.destroy_process_group()
Out[6]:
Console
Process group setup functions defined.
In real multi-GPU training, use 'nccl' backend and launch with torchrun.

Setting up the distributed process group is the first step in any DDP training script. The process group coordinates all communication between ranks. The gloo backend used here works on CPU and is suitable for demonstration; real GPU training always uses the nccl backend. The init_process_group call blocks until all processes in the group have called it, making sure that communication infrastructure is fully initialized before training begins.

Defining a Model and Wrapping with DDP

In[7]:
Code
import torch.nn as nn


class SimpleTransformerBlock(nn.Module):
    """A minimal transformer block for illustration."""

    def __init__(self, d_model=256, nhead=4, dim_feedforward=512):
        super().__init__()
        self.attn = nn.MultiheadAttention(d_model, nhead, batch_first=True)
        self.ff = nn.Sequential(
            nn.Linear(d_model, dim_feedforward),
            nn.ReLU(),
            nn.Linear(dim_feedforward, d_model),
        )
        self.norm1 = nn.LayerNorm(d_model)
        self.norm2 = nn.LayerNorm(d_model)

    def forward(self, x):
        attn_out, _ = self.attn(x, x, x)
        x = self.norm1(x + attn_out)
        x = self.norm2(x + self.ff(x))
        return x


class SimpleModel(nn.Module):
    def __init__(self, vocab_size=1000, d_model=256, num_layers=2):
        super().__init__()
        self.embed = nn.Embedding(vocab_size, d_model)
        self.layers = nn.ModuleList(
            [SimpleTransformerBlock(d_model) for _ in range(num_layers)]
        )
        self.head = nn.Linear(d_model, vocab_size)

    def forward(self, x):
        x = self.embed(x)
        for layer in self.layers:
            x = layer(x)
        return self.head(x)


# Count parameters
model = SimpleModel()
num_params = sum(p.numel() for p in model.parameters())
param_size_mb = sum(
    p.numel() * p.element_size() for p in model.parameters()
) / (1024**2)
Out[8]:
Console
Model parameters: 1,567,208
Model size (fp32): 6.0 MB
In DDP with N GPUs, each GPU holds 6.0 MB for model parameters.
Gradients require the same memory: 6.0 MB per GPU.
In[9]:
Code
from torch.nn.parallel import DistributedDataParallel as DDP

# Wrapping a model with DDP is a single line change
# device = torch.device("cuda:0")  # In real training
# model = model.to(device)
# ddp_model = DDP(model, device_ids=[local_rank])

# For demonstration (CPU / single-process):
setup_distributed()
ddp_model = DDP(model)
Out[10]:
Console
Model wrapped with DDP.
DDP hooks are now registered on all parameters.
Gradient all-reduce will be triggered automatically during backward().

Wrapping a model with DDP is the only model-related change needed in the training code. DDP registers backward hooks on all parameters automatically when the wrapper is created. When a gradient is computed during backpropagation, the hook fires and places the gradient into the appropriate bucket. When the bucket is full (or when the last parameter in the bucket gets its gradient), DDP launches the asynchronous NCCL all-reduce operation on that bucket.

One important detail: when you first construct the DDP wrapper, it broadcasts parameters from rank 0 to all other ranks. This ensures that all replicas start from identical initial parameters, even if rank 0 loaded a checkpoint and other ranks initialized randomly. After this initial broadcast, the invariant holds: all ranks have identical parameters at the start of every step.

The DistributedSampler

In[11]:
Code
import torch
from torch.utils.data import DataLoader, TensorDataset
from torch.utils.data.distributed import DistributedSampler

# Create synthetic training data
num_samples = 1024
seq_len = 32
vocab_size = 1000

X = torch.randint(0, vocab_size, (num_samples, seq_len))
y = torch.randint(0, vocab_size, (num_samples, seq_len))
dataset = TensorDataset(X, y)

# DistributedSampler ensures each rank processes a non-overlapping shard
# world_size is the number of GPUs; rank is the current GPU's index
sampler = DistributedSampler(
    dataset,
    num_replicas=1,  # world_size in real training
    rank=0,  # local_rank in real training
    shuffle=True,
    seed=42,
)

dataloader = DataLoader(dataset, batch_size=32, sampler=sampler)

# Count batches per rank
num_batches = len(dataloader)
samples_per_rank = len(sampler)
Out[12]:
Console
Total dataset size: 1,024 samples
Samples assigned to this rank: 1,024
Batches per rank per epoch: 32
With N GPUs, effective global batch size = 32 x N per step.

The DistributedSampler is critical for correctness. It partitions the dataset indices across ranks, making sure that each rank sees a different non-overlapping subset of samples each epoch. Without it, every rank would process the same samples, and you would be performing redundant computation instead of parallel computation. The effective batch size would remain bb rather than N×bN \times b, and you would get no speedup from additional GPUs.

The sampler also handles an important edge case: when the dataset size is not evenly divisible by the world size, some ranks would get slightly more samples than others. By default, DistributedSampler pads the dataset with replicated samples to ensure equal-length epochs across all ranks. This prevents DDP from hanging when one rank finishes its epoch and calls all-reduce while another rank still has samples to process.

The shuffle=True parameter causes the sampler to reshuffle the assignment of samples to ranks at the start of each epoch, using a shared seed based on the epoch number. This is why you must call sampler.set_epoch(epoch) before iterating the DataLoader each epoch: it updates the epoch counter used to compute the shuffled index assignments.

Training Loop with DDP

In[13]:
Code
import time

optimizer = torch.optim.AdamW(ddp_model.parameters(), lr=1e-3)
criterion = nn.CrossEntropyLoss()

# Simulate one epoch
ddp_model.train()
total_loss = 0.0
num_steps = 0
step_times = []

sampler.set_epoch(0)  # Must be called each epoch for correct shuffling

for batch_X, batch_y in dataloader:
    t0 = time.perf_counter()

    optimizer.zero_grad()

    # Forward pass: local computation only
    logits = ddp_model(batch_X)  # (B, T, vocab_size)
    loss = criterion(logits.view(-1, vocab_size), batch_y.view(-1))

    # Backward pass: DDP automatically all-reduces gradients
    loss.backward()

    # Gradient clipping (standard practice)
    nn.utils.clip_grad_norm_(ddp_model.parameters(), max_norm=1.0)

    # Optimizer step: identical update on all ranks
    optimizer.step()

    step_time = time.perf_counter() - t0
    step_times.append(step_time)
    total_loss += loss.item()
    num_steps += 1

avg_loss = total_loss / num_steps
avg_step_ms = np.mean(step_times) * 1000
Out[14]:
Console
Training steps completed: 32
Average loss: 7.0233
Average step time: 26.6 ms (CPU, single-process simulation)

In real DDP training:
  - backward() triggers async all-reduce on each gradient bucket
  - Communication overlaps with computation for earlier layers
  - All ranks apply identical optimizer updates
In[15]:
Code
# Cleanup
cleanup_distributed()

The training loop with DDP looks almost identical to single-GPU training. The only changes are wrapping the model with DDP, using DistributedSampler, and calling sampler.set_epoch(epoch) at the start of each epoch. The all-reduce happens automatically inside loss.backward(). This transparency is one of DDP's great strengths: existing training code requires minimal modification to become distributed.

One subtlety in the loop above is gradient clipping. The call to clip_grad_norm_ computes the global gradient norm across all parameters and clips it if it exceeds max_norm. Because all ranks have identical gradients after all-reduce, they compute identical norms and apply identical clipping. This is correct behavior. However, if you call clip_grad_norm_ before the all-reduce completes (which could happen in unusual training loop configurations), different ranks would clip based on their local gradient norms and diverge. DDP ensures all-reduce completes before the optimizer step, so standard gradient clipping is safe.

Key Parameters

The key parameters when setting up DDP are:

  • backend: Communication backend. Use "nccl" for GPU training (fastest, CUDA-aware), "gloo" for CPU training or debugging.
  • bucket_cap_mb: Maximum size of gradient buckets in megabytes (default 25). Larger buckets reduce communication launch overhead but increase memory usage and per-bucket latency before the first all-reduce starts.
  • find_unused_parameters: Set to True if some parameters may not receive gradients during every forward pass, for example in models with conditional routing or dynamic computation graphs. This adds overhead by scanning for unused parameters after each forward pass. Leave it False when all parameters are used in every forward pass, which is the common case for standard transformers.
  • gradient_as_bucket_view: When True, gradients are stored directly in the communication buffer rather than being copied into it. This saves a memory copy per bucket and reduces peak memory usage. It is recommended for memory-efficient training, though it slightly constrains how you can modify gradients between the backward pass and the optimizer step.
  • static_graph: When True, DDP assumes the computational graph is identical every step. This enables further optimizations including more aggressive bucket overlap and elimination of the find-unused-parameters scan. Set it to True for standard transformer models where the graph does not change.
  • device_ids: List of CUDA device IDs to use. Typically [local_rank] in multi-GPU setups where each process owns exactly one GPU.

Gradient Synchronization in Depth

The gradient synchronization step is where most of the complexity in DDP lives. Let us examine several aspects that matter in practice.

When Does Synchronization Happen?

DDP synchronizes gradients during the backward pass, not after it. This is achieved through autograd hooks that DDP registers on each parameter's .grad_fn. When the autograd engine finishes computing the gradient for a parameter, the hook fires immediately. DDP accumulates gradients into the appropriate bucket, and when the bucket is full (or when the last parameter in the bucket gets its gradient), DDP launches an asynchronous NCCL all-reduce operation on that bucket.

The hooks fire in the order that autograd computes gradients, which is roughly reverse layer order for a standard sequential model. This means the gradients for output layers (which backpropagate first) are placed into buckets and communicated first, while the gradients for early layers (embeddings, first transformer blocks) are computed and communicated last. By the time the backward pass finishes, many all-reduce operations have already completed.

This design has two important benefits. First, compute and communication run simultaneously: while the CPU continues the backward pass for earlier layers, the GPU's communication unit handles gradient transfers for later layers. Second, the critical path is shortened: the wall-clock time for a training step approaches max(compute_time, communication_time) rather than their sum, as long as the communication can be fully overlapped.

The practical degree of overlap depends on the model architecture. Deep models with many layers provide more opportunity for overlap than shallow models. Attention-heavy layers with long compute times provide more hiding of communication than lightweight MLP layers. For large transformer models (where each layer has hundreds of millions of parameters), the overlap is substantial and gradient bucketing provides large speedups.

No-Sync Context Manager

Sometimes you want to accumulate gradients over multiple mini-batches before synchronizing. This is called gradient accumulation, and it is used when the target global batch size is larger than what fits in GPU memory in a single batch. The idea is to run kk forward-backward passes and only call optimizer.step() after kk passes, effectively computing a gradient over kk times as many samples per optimizer step.

With DDP, calling backward() triggers an all-reduce after every pass, even when you do not want to synchronize yet. This wastes communication bandwidth: you are averaging gradients kk times per optimizer step but only need to average them once. DDP provides the no_sync() context manager to suppress gradient synchronization:

In[16]:
Code
# Gradient accumulation with DDP
accumulation_steps = 4
optimizer.zero_grad()

ddp_model.train()
sampler.set_epoch(1)

for step, (batch_X, batch_y) in enumerate(dataloader):
    # Only synchronize on the last accumulation step
    is_last_accumulation_step = (step + 1) % accumulation_steps == 0

    if is_last_accumulation_step:
        # Normal backward: triggers all-reduce
        logits = ddp_model(batch_X)
        loss = criterion(logits.view(-1, vocab_size), batch_y.view(-1))
        loss = loss / accumulation_steps
        loss.backward()
        optimizer.step()
        optimizer.zero_grad()
    else:
        # Suppressed backward: no all-reduce until no_sync exits
        with ddp_model.no_sync():
            logits = ddp_model(batch_X)
            loss = criterion(logits.view(-1, vocab_size), batch_y.view(-1))
            loss = loss / accumulation_steps
            loss.backward()

    if step >= 7:  # Run just a few steps for demonstration
        break

num_sync_steps = (step + 1) // accumulation_steps
Out[17]:
Console
Ran 8 forward-backward passes.
Triggered all-reduce 2 times.
Communication savings vs naive: 6 all-reduces skipped.

Gradient accumulation allows larger effective batch sizes
without increasing per-step communication frequency.

Gradient accumulation is an important practical technique because many language model training recipes target specific global batch sizes (for example, 4 million tokens per step) that far exceed what fits in GPU memory in a single forward pass. By accumulating over kk steps with no_sync(), you achieve the target batch size while only communicating once. The communication overhead is the same as without accumulation, but the effective batch size is kk times larger.

One thing to be careful about: the loss / accumulation_steps scaling in the code above. Because gradients accumulate additively across the kk backward passes, you need to scale each loss by 1/k1/k so that the accumulated gradient equals the gradient of the mean loss over the full accumulated batch. Without this scaling, the accumulated gradient would be kk times the mean gradient, which would require reducing the learning rate by kk to compensate.

Gradient Averaging vs. Gradient Summing

DDP averages gradients by default, dividing the sum by NN (the number of ranks). This normalization ensures that the gradient scale is the same regardless of how many GPUs you use, which makes learning rate tuning independent of the degree of parallelism.

There is an important subtlety here. The per-rank loss is already averaged over the local batch:

Lr=1b∑x∈BrLx(θ)\mathcal{L}_r = \frac{1}{b} \sum_{x \in \mathcal{B}_r} \mathcal{L}_x(\theta)

where b=∣Br∣b = |\mathcal{B}_r| is the local batch size. The gradient of this loss is:

∇Lr=1b∑x∈Br∇Lx(θ)\nabla \mathcal{L}_r = \frac{1}{b} \sum_{x \in \mathcal{B}_r} \nabla \mathcal{L}_x(\theta)

After all-reduce averaging over NN ranks, the synchronized gradient is:

∇sync=1N∑r=0N−1∇Lr=1N∑r=0N−11b∑x∈Br∇Lx(θ)=1Nb∑r=0N−1∑x∈Br∇Lx(θ)\begin{aligned} \nabla_{\text{sync}} &= \frac{1}{N} \sum_{r=0}^{N-1} \nabla \mathcal{L}_r \\ &= \frac{1}{N} \sum_{r=0}^{N-1} \frac{1}{b} \sum_{x \in \mathcal{B}_r} \nabla \mathcal{L}_x(\theta) \\ &= \frac{1}{Nb} \sum_{r=0}^{N-1} \sum_{x \in \mathcal{B}_r} \nabla \mathcal{L}_x(\theta) \end{aligned}

where:

  • ∇sync\nabla_{\text{sync}} is the final synchronized gradient used for the optimizer step
  • bb is the per-rank mini-batch size
  • NN is the number of ranks (GPUs)

This expression is exactly the gradient of the mean loss over the global batch of size NbNb. The synchronized gradient is the true full-batch gradient for the combined mini-batch across all ranks, divided by the global batch size. This is the correct gradient for the loss function you are optimizing, which is the mean loss per token (or per example) over the global batch.

Scaling Behaviour

Understanding how DDP scales is essential for planning distributed training runs. Ideal scaling and real-world scaling differ in important ways, and knowing these differences prevents unpleasant surprises when you provision hardware.

Ideal Linear Scaling

In the ideal case, adding NN GPUs gives you an NN-fold speedup. With batch size bb per GPU, each training step processes N×bN \times b total samples. Since one step on NN GPUs takes the same wall-clock time as one step on 1 GPU (assuming communication is free), you process NN times as many samples per unit time. Training runs NN times faster.

This ideal scaling holds when:

  1. Communication time is negligible compared to compute time
  2. The workload divides evenly among GPUs (no load imbalance)
  3. All GPUs complete their local computation at the same time (no stragglers)
  4. Data loading keeps up with GPU throughput

For large models with many parameters, condition 1 often approximately holds because compute time scales with model size while communication volume also scales with model size but at a lower constant factor (ring all-reduce transfers only 2∣g∣2|\mathbf{g}| bytes per rank regardless of NN). A 7B model has 14 GB of BF16 gradients but requires several hundred milliseconds of compute per step, so the communication-to-compute ratio can be small on fast interconnects.

Amdahl's Law and Communication Overhead

In practice, some fraction of each training step is inherently serial or requires synchronization. Amdahl's Law captures this fundamental limit: if fraction pp of the work can be parallelized and fraction 1−p1-p is serial, the maximum speedup from NN processors is:

S(N)=1(1−p)+pNS(N) = \frac{1}{(1-p) + \frac{p}{N}}

where:

  • S(N)S(N) is the speedup with NN GPUs relative to 1 GPU
  • pp is the fraction of total time that scales with more GPUs (the parallelizable part)
  • 1−p1-p is the serial fraction (overhead from synchronization, data loading, communication, and other fixed costs)

As N→∞N \to \infty, the speedup approaches 11−p\frac{1}{1-p}. Even with a serial fraction of just 5%, the maximum possible speedup is 20x regardless of how many GPUs you add. At 10% serial fraction, the ceiling drops to 10x. These limits are strict: no amount of engineering can overcome them without reducing the serial fraction.

For DDP, the main non-parallelizable overhead is the all-reduce communication that cannot be fully overlapped with compute. The communication time for ring all-reduce is approximately:

tcomm≈2⋅N−1N⋅∣g∣Blinkt_{\text{comm}} \approx 2 \cdot \frac{N-1}{N} \cdot \frac{|\mathbf{g}|}{B_{\text{link}}}

where:

  • ∣g∣|\mathbf{g}| is the gradient tensor size in bytes
  • BlinkB_{\text{link}} is the available link bandwidth per GPU in bytes per second (this is bidirectional bandwidth; NCCL can use both send and receive channels simultaneously)
  • The factor approaches 2 as NN grows large (both reduce-scatter and all-gather phases each transfer approximately ∣g∣|\mathbf{g}| bytes per rank)

This communication time is independent of NN for large NN, because ring all-reduce is bandwidth-optimal. The implication is that DDP scales very well within a single node (NVLink bandwidth is high, so communication is fast) but scales less well across nodes (InfiniBand bandwidth is lower). Each additional node adds another cross-node communication hop, and if the cross-node bandwidth is already saturated, scaling efficiency degrades.

The Linear Scaling Rule

When scaling DDP to more GPUs, you should also scale the global learning rate. The standard practice, established by Goyal et al. (2017) in their ImageNet training paper and confirmed for language models by many subsequent works, is the linear scaling rule:

αN=α1×N\alpha_N = \alpha_1 \times N

where:

  • α1\alpha_1 is the learning rate for 1-GPU training with local batch size bb
  • αN\alpha_N is the learning rate for NN-GPU training with global batch size NbNb

The intuition is this: with NN times as many samples per update step, each gradient estimate is NN times less noisy (by the central limit theorem, the variance of a mean over NN samples is 1/N1/N times the variance of a single sample). Lower gradient noise means each update step is more reliable, so you can afford to take NN times larger steps. Taking larger steps means you cover the same loss landscape in fewer steps, maintaining equivalent training dynamics.

The linear scaling rule works well for batch sizes up to a few thousand examples. For very large batch sizes (beyond 8,000 to 16,000 for typical language model training), learning rate scaling alone is insufficient. At these scales, the gradient noise is so low that each update consumes nearly all of the useful curvature information available in the loss landscape in a single step, and you enter the "large batch regime" where additional batch size no longer reduces the number of steps needed proportionally. Research by McCandlish et al. (2018) on the "critical batch size" formalized this regime.

Learning rate warmup is also essential when scaling to large batch sizes. Rather than jumping immediately to the scaled learning rate αN\alpha_N, you linearly increase the learning rate from a small initial value (or zero) to αN\alpha_N over the first few thousand training steps. This warmup period prevents instability at the start of training when the optimizer's first and second moment estimates (in Adam/AdamW) have not yet converged and the model is far from any sensible initialization.

Code: Scaling Analysis

Let us simulate the scaling behaviour of DDP and examine the impact of communication overhead analytically.

In[18]:
Code
# Model parameters affect gradient size
param_configs = {
    "125M (GPT-2 small)": 125e6,
    "1.3B (GPT-2 XL equiv)": 1.3e9,
    "7B (Llama-7B equiv)": 7e9,
}

# Hardware parameters
nvlink_bandwidth_gb = 600  # NVLink bandwidth per GPU in GB/s (intranode)
infiniband_bandwidth_gb = (
    25  # InfiniBand effective bandwidth per GPU in GB/s (internode)
)
nvlink_step_latency_ms = 0.02
infiniband_step_latency_ms = 0.05
bytes_per_param_bf16 = 2  # BF16 gradients

# Illustrative compute times chosen to expose the compute/communication balance
compute_times_ms = {
    "125M (GPT-2 small)": 3.0,
    "1.3B (GPT-2 XL equiv)": 50.0,
    "7B (Llama-7B equiv)": 400.0,
}


def ring_allreduce_time_ms(num_params, bandwidth_gb_per_s, N, step_latency_ms):
    """Estimate ring all-reduce time in milliseconds."""
    gradient_bytes = num_params * bytes_per_param_bf16
    bandwidth_bytes_per_ms = bandwidth_gb_per_s * 1e9 / 1e3
    # Ring all-reduce: send+receive ~ 2 * (N-1)/N * gradient_bytes
    comm_bytes = 2 * ((N - 1) / N) * gradient_bytes
    bandwidth_time_ms = comm_bytes / bandwidth_bytes_per_ms
    latency_time_ms = 2 * (N - 1) * step_latency_ms
    return bandwidth_time_ms + latency_time_ms


def compute_speedup(
    N_values, num_params, compute_ms, bandwidth_gb, step_latency_ms
):
    """Compute actual vs ideal speedup for a range of GPU counts."""
    speedups = []
    efficiencies = []
    for N in N_values:
        comm_ms = ring_allreduce_time_ms(
            num_params, bandwidth_gb, N, step_latency_ms
        )
        # Step time = compute (overlapped with comm) = max(compute, comm)
        # Plus a small serial fraction for sync, data loading, etc.
        serial_overhead_ms = 0.5  # fixed overhead per step
        step_time = max(compute_ms, comm_ms) + serial_overhead_ms
        single_gpu_step_ms = compute_ms + serial_overhead_ms
        speedup = (
            single_gpu_step_ms * N / step_time
        )  # effective throughput gain
        speedups.append(speedup)
        efficiencies.append(speedup / N * 100)
    return speedups, efficiencies


N_values = [1, 2, 4, 8, 16, 32, 64, 128]
Out[19]:
Console
Configuration summary:
Model                                Params   Grad size (BF16)
--------------------------------------------------------------
125M (GPT-2 small)                125000000             238 MB
1.3B (GPT-2 XL equiv)            1300000000            2480 MB
7B (Llama-7B equiv)              7000000000           13351 MB
In[20]:
Code
# Compute scaling curves for each model and bandwidth scenario
results = {}
for model_name, num_params in param_configs.items():
    compute_ms = compute_times_ms[model_name]
    speedups_nv, efficiencies_nv = compute_speedup(
        N_values,
        num_params,
        compute_ms,
        nvlink_bandwidth_gb,
        nvlink_step_latency_ms,
    )
    speedups_ib, efficiencies_ib = compute_speedup(
        N_values,
        num_params,
        compute_ms,
        infiniband_bandwidth_gb,
        infiniband_step_latency_ms,
    )
    results[model_name] = {
        "speedup_nvlink": speedups_nv,
        "speedup_infiniband": speedups_ib,
        "efficiency_nvlink": efficiencies_nv,
        "efficiency_infiniband": efficiencies_ib,
    }

# Compute ideal speedup
ideal_speedup = N_values
Out[21]:
Console
Scaling efficiency for 7B (Llama-7B equiv):
  GPUs   NVLink speedup  NVLink eff%   InfiniBand speedup    IB eff%
----------------------------------------------------------------------
     1              1.0        100.0                  1.0      100.0
     2              2.0        100.0                  1.4       71.4
     4              4.0        100.0                  1.9       47.6
     8              8.0        100.0                  3.3       40.8
    16             16.0        100.0                  6.1       38.1
    32             32.0        100.0                 11.8       36.8
    64             64.0        100.0                 23.1       36.1
   128            128.0        100.0                 45.6       35.6
Out[22]:
Visualization
Line chart showing DDP speedup vs GPU count for 125M, 1.3B, and 7B models on NVLink.
DDP scaling speedup vs number of GPUs for three model sizes using NVLink (intranode, 600 GB/s). Larger models scale more efficiently because compute time dominates communication time, keeping scaling efficiency high. The 7B model remains near linear through 128 GPUs in this illustrative model, while the 125M curve falls increasingly below ideal beyond 32 GPUs as collective latency overtakes its shorter compute time.
Line chart showing DDP speedup vs GPU count for 125M, 1.3B, and 7B models on InfiniBand.
DDP scaling speedup vs number of GPUs using InfiniBand (internode, 25 GB/s effective). Communication time becomes a bottleneck for all but the largest models, particularly beyond 32 GPUs. The 125M model scales poorly because its small compute time is overwhelmed by communication overhead at even modest GPU counts.

The speedup curves reveal a fundamental truth about data parallelism: larger models scale more efficiently than smaller ones. For a 7B parameter model on NVLink, compute time is so much larger than communication time that DDP achieves near-linear scaling up to 64 GPUs. For a 125M model on InfiniBand across nodes, the relatively short compute time is quickly dominated by communication overhead, and scaling efficiency collapses beyond 8 to 16 GPUs.

This has direct practical implications. If you are training a large language model, data parallelism is an excellent fit because the models are large enough that compute always dominates communication. If you are training a smaller model or performing frequent experimentation with small configurations, you may hit communication bottlenecks sooner and need to consider whether the parallelism overhead is worth the speedup.

Out[23]:
Visualization
Heatmap of DDP scaling efficiency percent on NVLink for 6 model sizes and 8 GPU counts.
DDP scaling efficiency (percentage of ideal speedup achieved) as a function of GPU count and model size on NVLink (intranode). Efficiency remains above 85% for larger models at all GPU counts shown, because compute time dwarfs communication time. Smaller models lose efficiency faster as GPU count grows, since their lower compute-to-communication ratio means all-reduce takes a growing fraction of each step.
Heatmap of DDP scaling efficiency percent on InfiniBand for 6 model sizes and 8 GPU counts.
DDP scaling efficiency on InfiniBand (internode, 25 GB/s). Cross-node communication is 24x slower than NVLink, causing efficiency to drop sharply for small and medium models beyond 8 GPUs. Only the largest models (7B and above) maintain useful efficiency at 32 or more GPUs on InfiniBand, making within-node NVLink training strongly preferred when possible.

Communication-Compute Overlap Efficiency

In[24]:
Code
# Simulate the effect of gradient bucketing on step time
# With perfect overlap: step_time = max(compute, comm)
# With no overlap: step_time = compute + comm
# Reality is between these extremes


def step_time_analysis(compute_ms, comm_ms, overlap_fraction=0.85):
    """
    Estimate step time with partial overlap.
    overlap_fraction: fraction of communication that overlaps with computation.
    """
    comm_serial = comm_ms * (1 - overlap_fraction)
    comm_overlapped = comm_ms * overlap_fraction

    # Overlapped communication hides within compute time
    effective_step = compute_ms + comm_serial
    return effective_step


model_name = "7B (Llama-7B equiv)"
num_params = param_configs[model_name]
compute_ms = compute_times_ms[model_name]

print(f"Analysis for {model_name}:")
print(f"Compute time per step: {compute_ms:.1f} ms")
print()

for N in [1, 2, 4, 8, 16, 32]:
    comm_nv = ring_allreduce_time_ms(
        num_params, nvlink_bandwidth_gb, N, nvlink_step_latency_ms
    )
    comm_ib = ring_allreduce_time_ms(
        num_params, infiniband_bandwidth_gb, N, infiniband_step_latency_ms
    )
    step_nv = step_time_analysis(compute_ms, comm_nv, overlap_fraction=0.85)
    step_ib = step_time_analysis(compute_ms, comm_ib, overlap_fraction=0.85)
    throughput_nv = N * compute_ms / step_nv  # relative throughput
    throughput_ib = N * compute_ms / step_ib
    print(
        f"  N={N:3d}  NVLink comm={comm_nv:5.1f} ms  IB comm={comm_ib:5.1f} ms  "
        f"  NVLink eff={throughput_nv:.1f}x  IB eff={throughput_ib:.1f}x"
    )
Out[25]:
Console
Analysis for 7B (Llama-7B equiv):
Compute time per step: 400.0 ms

   N  NVLink comm    IB comm   NVLink eff     IB eff
----------------------------------------------------
   1        0.0 ms      0.0 ms        1.0x      1.0x
   2       23.4 ms    560.1 ms        2.0x      1.7x
   4       35.1 ms    840.3 ms        3.9x      3.0x
   8       41.1 ms    980.7 ms        7.9x      5.8x
  16       44.4 ms   1051.5 ms       15.7x     11.5x
  32       46.4 ms   1088.1 ms       31.5x     22.7x

The communication time for ring all-reduce on a 7B model is substantial. On NVLink, the 14 GB of BF16 gradients can be transferred in roughly 45 ms at 600 GB/s. On InfiniBand at 25 GB/s effective bandwidth, the same communication takes over a second. With 250 ms of compute time, NVLink communication is safely within the compute window and can be largely overlapped. InfiniBand communication is four times the compute time, severely limiting scaling efficiency beyond a single node. This quantitative comparison explains why multi-node training at scale typically combines DDP with ZeRO optimization: ZeRO reduces the volume of gradient traffic that must cross the slow inter-node fabric.

Beyond DDP: Fully Sharded Data Parallel

Pure DDP replicates the complete model on every GPU. This is simple and efficient when the model fits in GPU memory, but as models grow beyond the memory capacity of a single GPU (80 GB for A100/H100), pure DDP is impossible without additional techniques.

Fully Sharded Data Parallel (FSDP), introduced by Facebook AI Research and available in PyTorch as torch.distributed.fsdp.FullyShardedDataParallel, extends DDP by sharding the data, model parameters, gradients, and optimizer states across GPUs.

How FSDP Works

In FSDP, each GPU holds only 1/N1/N of the parameters at any given time. Before the forward pass through each layer, FSDP gathers the full parameter tensor for that layer from all GPUs (an all-gather operation). After the forward pass for that layer, FSDP discards the gathered parameters, retaining only its own shard. During the backward pass, FSDP re-gathers parameters on demand, computes gradients, reduces those gradients across GPUs (a reduce-scatter), and again discards the gathered parameters.

The memory savings are dramatic. For a 7B model with AdamW in mixed precision:

  • Pure DDP: 16 bytes per parameter times 7 billion = 112 GB per GPU
  • FSDP (with optimizer state sharding): approximately 16 bytes per parameter divided by NN GPUs = 112/N GB per GPU

With N=8N = 8 GPUs, FSDP reduces per-GPU memory from 112 GB to 14 GB, enabling training that is otherwise impossible.

The cost of FSDP is increased communication volume. FSDP performs additional all-gather operations for parameter reconstruction before each layer's forward and backward pass. This roughly doubles the communication volume compared to pure DDP. However, FSDP can overlap these all-gather operations with the computation of the preceding layer, similar to how DDP overlaps gradient all-reduce with the backward pass.

FSDP vs. DDP: When to Use Which

The choice between DDP and FSDP depends primarily on model size relative to GPU memory:

  • If your model (plus gradients and optimizer states) fits on a single GPU, use DDP. It is simpler, has lower communication overhead, and is easier to debug.
  • If your model does not fit on a single GPU but you need data parallelism, use FSDP. It achieves the same data-parallel scaling as DDP while reducing per-GPU memory by a factor of NN.
  • For the largest models (hundreds of billions of parameters), even FSDP may be insufficient without also combining tensor parallelism or pipeline parallelism, which we cover in subsequent chapters.

The ZeRO optimizer (covered in the ZeRO chapter) implements the same sharding strategy as FSDP through the DeepSpeed library. ZeRO Stage 1 shards only optimizer states, Stage 2 adds gradient sharding, and Stage 3 adds parameter sharding. Stage 3 is roughly equivalent to FSDP in terms of memory savings and communication volume.

Limitations and Trade-offs

Despite its simplicity and effectiveness, pure data parallelism has fundamental constraints that become relevant as models and cluster sizes grow. Understanding these limitations shapes how you combine DDP with other distributed training strategies.

Memory Replication

Data parallelism's most significant cost is memory: every GPU must hold a complete copy of the model, gradients, and optimizer states. For a model with PP parameters trained with AdamW in mixed precision, the memory cost per GPU is approximately:

Memory per GPU=P⋅(2+2+4+8)=16P bytes\text{Memory per GPU} = P \cdot (2 + 2 + 4 + 8) = 16P \text{ bytes}

where the terms represent: 2 bytes for BF16 parameters, 2 bytes for BF16 gradients, 4 bytes for FP32 master weights (used by the optimizer for numerical stability), and 8 bytes for Adam's first and second moment estimates stored in FP32 (4 bytes each).

For a 7B parameter model, this is roughly 112 GB per GPU, far exceeding the 80 GB capacity of an A100 or H100. Pure data parallelism cannot train such models without additional techniques. ZeRO optimization (Chapter 6 of this part) addresses this by partitioning the optimizer states and gradients across GPUs. FSDP goes further and also partitions the parameters. Tensor parallelism and pipeline parallelism (Chapters 4 and 5) partition the model computation itself across GPUs.

This is why real-world large-scale training combines multiple parallelism strategies. Data parallelism is typically the outermost dimension: you have DD data-parallel replicas, and each replica itself uses tensor parallelism across TT GPUs within a node. The total GPU count is D×TD \times T. The gradient all-reduce in the DDP dimension only needs to communicate between the DD data-parallel groups, not across all D×TD \times T GPUs.

Out[26]:
Visualization
Stacked bar chart showing per-GPU memory breakdown for DDP across three model sizes.
Per-GPU memory breakdown for DDP training with AdamW in mixed precision, shown for three model sizes (125M, 1.3B, 7B). Each GPU holds BF16 parameters and gradients (2 bytes each), FP32 master weights (4 bytes), and Adam first and second moment estimates (8 bytes total), summing to 16 bytes per parameter. The dashed line marks the 80 GB capacity of an A100 or H100 GPU. The 7B model exceeds single-GPU capacity, requiring ZeRO optimization or model parallelism to fit.

Synchronization Stalls and Stragglers

DDP requires all ranks to reach the barrier at the end of each step before any rank can proceed to the next step. This is enforced by the all-reduce: every rank must contribute its gradients, and no rank receives the averaged result until all ranks have submitted. If one GPU is slower than the others (due to hardware variation, variable batch complexities, thermal throttling, or OS preemption), all other GPUs must wait at the barrier. This straggler effect directly reduces effective GPU utilization.

The straggler problem grows worse as the number of GPUs increases. With N=8N = 8 GPUs, the slowest GPU on average is somewhat faster than with N=512N = 512 GPUs, where the probability of having a significantly slow GPU in any given step is much higher. In large clusters, per-step variance in GPU timing is substantial, and straggler management becomes a first-class concern.

Practical straggler mitigation strategies include:

  • Using homogeneous hardware (all GPUs the same model and generation)
  • Setting appropriate timeouts on collective operations to detect and recover from hangs
  • Using padding and masking for variable-length sequences rather than variable-length batches (which would cause load imbalance across ranks)
  • Monitoring per-rank step times to identify consistently slow ranks (often indicating hardware faults)

PyTorch's distributed training infrastructure includes a watchdog mechanism that detects collective operation timeouts and raises informative errors, preventing silent hangs that would otherwise waste compute indefinitely.

The Global Batch Size Problem

Scaling to more GPUs means a larger global batch size, which can negatively affect model quality. Research has shown that there is a sweet spot for batch size in language model training. Below this "critical batch size," doubling the batch size roughly halves the number of training steps needed to reach a given loss level, so the total compute is constant. Above the critical batch size, further increases in batch size no longer produce proportional reductions in step count: you need more total compute (more tokens seen) to reach the same quality.

The critical batch size for language models is roughly related to the signal-to-noise ratio of the gradient. As the batch size increases, gradient variance decreases (noise goes down), but the useful signal (the true gradient direction) also becomes "used up" faster per update. The optimal regime balances these effects. Empirically, the critical batch size for large language model training is typically in the range of a few hundred thousand to a few million tokens per step.

When you add GPUs and correspondingly increase the global batch size, you may exceed this critical batch size. Training with an above-critical batch size is not wrong: the model will still converge, and you are not wasting compute in an absolute sense. But you may need to train for more total tokens than you would have with a smaller batch size to reach the same final quality. The linear scaling rule adjusts for this in the optimizer step size, but it cannot eliminate the fundamental relationship between batch size and gradient quality.

Grasping this trade-off is essential for practical distributed training decisions. Adding GPUs is not free in terms of model quality per training token. The optimal configuration depends on your hardware budget, time budget, and quality requirements. For many practical training scenarios, a combination of data parallelism, gradient accumulation, and careful batch size tuning achieves the best balance.

Data Loading Bottlenecks

With many GPUs processing data rapidly, the data loading pipeline can become a bottleneck. Each rank needs its data on time for every step. With insufficient I/O throughput, GPUs stall waiting for data rather than doing computation. For large transformer models that process data quickly, this is a real risk, particularly when training on slow network filesystems or when tokenization is performed on-the-fly.

Solutions include:

  • Multiple DataLoader workers: Use num_workers=8 or more in DataLoader to parallelize data fetching and preprocessing across CPU threads.
  • Pinned memory: Use pin_memory=True in DataLoader to store batches in page-locked host memory, enabling faster CPU-to-GPU memory transfers via DMA.
  • Pre-tokenized data: Store tokenized data in memory-mapped binary files (for example, NumPy .npy, HuggingFace datasets with Arrow format, or webdataset tar files) to avoid per-epoch tokenization overhead.
  • Streaming datasets: For very large datasets (trillions of tokens), stream data from disk rather than loading the whole dataset into RAM. Libraries like webdataset and HuggingFace's streaming datasets support this pattern.
  • Local SSD caching: Copy frequently accessed data from slow network storage to fast local NVMe SSDs on training nodes to reduce I/O latency.

Fault Tolerance

Pure DDP has no built-in fault tolerance. If any rank fails (due to a hardware error, preemption, or network partition), the training job fails completely. The all-reduce collective requires all participants; any absent rank causes all others to hang indefinitely at the barrier.

For short training runs on reliable hardware, this is acceptable. For multi-week training runs on thousands of GPUs, hardware failures are not rare events but expected occurrences. GPU failures, network transients, and node preemptions happen daily on large clusters. Without fault tolerance, a single failure wastes all progress since the last checkpoint.

Solutions include:

  • Frequent checkpointing: Save model and optimizer state every few hundred steps. Combined with elastic training frameworks that can restart on fewer GPUs, checkpoints bound the amount of wasted compute after a failure.
  • Elastic training: PyTorch's torchrun and frameworks like PyTorch Lightning support elastic training, where the number of workers can change dynamically and training resumes with a different number of GPUs after a failure.
  • Redundant training: Some large-scale training runs checkpoint to distributed storage systems (like GCS or S3) that are themselves fault-tolerant, making sure checkpoints survive even if the training nodes fail.

Summary

Data parallelism is the foundational strategy for distributed language model training. By maintaining a complete model copy on each GPU and processing different data shards in parallel, DDP achieves near-linear speedup when communication bandwidth is sufficient relative to compute time.

The key ideas to carry forward are:

  • DDP algorithm: Each rank computes local gradients on its data shard. All-reduce synchronizes gradients across all ranks. All ranks apply identical parameter updates. The model copies stay perfectly synchronized throughout training without any explicit parameter communication.
  • All-reduce: Ring all-reduce is the bandwidth-optimal primitive that makes DDP scalable. Each rank transfers roughly 2∣g∣2|\mathbf{g}| bytes regardless of the number of GPUs, distributing aggregation work evenly across the ring. NCCL implements topology-aware variants that exploit NVLink for intranode communication and InfiniBand for cross-node traffic.
  • Gradient bucketing: DDP overlaps communication with computation by launching all-reduce on completed gradient buckets during the backward pass, hiding most communication latency. The bucket size (default 25 MB) controls the tradeoff between communication launch overhead and overlap granularity.
  • Scaling limits: Larger models scale more efficiently because their high compute-to-communication ratio keeps GPUs busy. Cross-node training with InfiniBand scales worse than within-node NVLink due to roughly 24x bandwidth difference. The linear scaling rule adjusts learning rates for larger global batch sizes, but above the critical batch size, total token cost increases.
  • Memory cost: Pure data parallelism replicates the full model on every GPU, requiring 16 bytes per parameter for AdamW mixed precision training. For models exceeding GPU memory (above roughly 5B parameters on an A100), FSDP, ZeRO optimization, or model parallelism must supplement data parallelism.
  • Practical concerns: Data loading, straggler GPUs, batch size tuning, and fault tolerance all require attention in production distributed training setups.

The next chapters cover complementary strategies that work alongside data parallelism. Tensor parallelism (Chapter 4) splits individual layers across GPUs within a node, reducing per-GPU memory proportionally. Pipeline parallelism (Chapter 5) stages model layers across sets of GPUs, enabling even larger models. ZeRO optimization (Chapter 6) achieves memory reduction equivalent to model parallelism while retaining the simpler programming model of data parallelism. Together, these form the parallelism stack used by every large-scale language model training run today.

Quiz

Ready to test your understanding? Take this quick quiz to reinforce what you've learned about data parallelism, DDP, and gradient synchronization.

Data Parallelism Quiz

Question 1 of 80 of 8 completed
In Distributed Data Parallel (DDP), what does each GPU (rank) hold during training?

Comments

No comments yet. Be the first to share your thoughts!

Reference

Citation details

Cite or share this article.

BIBTEXAcademic
@misc{brenndoerfer2026dataparallelism, author = {Michael Brenndoerfer}, title = {Data Parallelism: DDP, Gradient Synchronization & All-Reduce}, year = {2026}, url = {https://mbrenndoerfer.com/writing/data-parallelism-ddp-gradient-synchronization-all-reduce}, organization = {mbrenndoerfer.com}, note = {Accessed: 2026-10-06} }
APAAcademic
Michael Brenndoerfer (2026). Data Parallelism: DDP, Gradient Synchronization & All-Reduce. Retrieved from https://mbrenndoerfer.com/writing/data-parallelism-ddp-gradient-synchronization-all-reduce
MLAAcademic
Michael Brenndoerfer. "Data Parallelism: DDP, Gradient Synchronization & All-Reduce." 2026. Web. October 6, 2026. <https://mbrenndoerfer.com/writing/data-parallelism-ddp-gradient-synchronization-all-reduce>.
CHICAGOAcademic
Michael Brenndoerfer. "Data Parallelism: DDP, Gradient Synchronization & All-Reduce." Accessed October 6, 2026. https://mbrenndoerfer.com/writing/data-parallelism-ddp-gradient-synchronization-all-reduce.
HARVARDAcademic
Michael Brenndoerfer (2026) 'Data Parallelism: DDP, Gradient Synchronization & All-Reduce'. Available at: https://mbrenndoerfer.com/writing/data-parallelism-ddp-gradient-synchronization-all-reduce (Accessed: October 6, 2026).
SimpleBasic
Michael Brenndoerfer (2026). Data Parallelism: DDP, Gradient Synchronization & All-Reduce. https://mbrenndoerfer.com/writing/data-parallelism-ddp-gradient-synchronization-all-reduce

About the author

Continue with the full handbook

This chapter is part of Language AI Handbook. Use the handbook page to browse the complete table of contents and continue reading in sequence.

Explore Language AI Handbook
Newsletter

Stay up to date

Get articles, book updates, and news delivered to your inbox.

No spam, unsubscribe anytime.

or

Join the community

Sign in to remove popups, track your reading progress, and join the discussion.