Every multi-GPU training job is a sequence of local compute punctuated by collective operations: calls in which every GPU in a group takes part, and each ends with data that depends on everyone else's. The data-parallel gradient sum is a collective; so are the parameter gathers inside FSDP, the activation exchange in tensor parallelism and the token shuffle in a mixture-of-experts layer. If you know which collective each strategy issues, how many bytes it moves and what bounds its speed, you can predict step time before you launch, read a profile when it is slow and choose a parallelism layout on purpose.

This article is the map. It defines each collective precisely, shows which training scheme uses which, builds a simple cost model, measures it with nccl-tests and PyTorch, works through an example, and lists the failure modes that hang or slow real jobs. How NCCL chooses ring or tree algorithms over a given topology is covered separately in NCCL collectives; here the focus is what the operations are and what your training code asks of them.

The collective vocabulary

Take n ranks (one process per GPU), each holding a buffer. The core collectives differ in what each rank has before and after. In the table, S is the size of the full tensor and the last column is the minimum data each rank must send in a bandwidth-optimal implementation.

CollectiveBefore (each rank)After (each rank)Bytes sent per rank
Broadcastroot has Sall have root's Sabout S (pipelined)
ReduceS eachroot has the sumabout S
All-reduceS eachall have the sum2(n-1)/n x S
Reduce-scatterS eachone reduced shard of S/n(n-1)/n x S
All-gatherone shard of S/nall n shards (S)(n-1)/n x S
All-to-alln blocks of S/nblock i from every rank(n-1)/n x S
Send / recvpoint to pointpeer has the datamessage size

Two identities are worth remembering. All-reduce is reduce-scatter followed by all-gather, which is why its traffic factor is twice theirs. And the factors approach a constant as n grows: each rank sends just under twice the tensor for an all-reduce no matter whether there are 8 or 8,000 ranks. That is what makes data parallelism scale in bandwidth terms. Latency is a different story, covered below.

All-reduce = reduce-scatter + all-gather (4 ranks, tensor split into 4 shards)rank 0: a0 b0 c0 d0full local gradientrank 1: a1 b1 c1 d1full local gradientrank 2: a2 b2 c2 d2full local gradientrank 3: a3 b3 c3 d3full local gradientreduce-scatter: each rank sends 3/4 of the tensor, keeps one summed shardrank 0: sum a1/4 of tensor, reducedrank 1: sum b1/4 of tensor, reducedrank 2: sum c1/4 of tensor, reducedrank 3: sum d1/4 of tensor, reducedall-gather: each rank sends its shard, receives the other 3rank 0: all sumsfull reduced tensorrank 1: all sumsfull reduced tensorrank 2: all sumsfull reduced tensorrank 3: all sumsfull reduced tensorFSDP stops after the first stage (keeps a shard); DDP runs both stages.Per rank, each stage moves (n-1)/n of the tensor, so all-reduce moves 2(n-1)/n.
Figure: the two stages of an all-reduce. Sharded data parallelism keeps the result of the first stage and runs the second stage on parameters instead of gradients.

Which training strategy issues which collective

Each parallelism strategy is defined by which collective it issues, on which tensor and how often. This mapping is the most useful thing to carry in your head when reading a profile.

StrategyCollectiveTensorWhen
Data parallel (DDP)All-reducegradientsonce per step, bucketed during backward
FSDP / ZeRO-3All-gatherparametersbefore each layer's forward and backward
FSDP / ZeRO-3Reduce-scattergradientsafter each layer's backward
ZeRO-1/2Reduce-scatter, then all-gathergrads, then updated paramsper step
Tensor parallel (Megatron style)All-reduceactivations, activation gradstwice per transformer block, each pass
Sequence parallel with TPAll-gather + reduce-scatteractivationsreplaces the TP all-reduces
Expert parallel (MoE)All-to-alltokensdispatch and combine, every MoE layer
Pipeline parallelSend / recvactivations, gradsbetween stages, per micro-batch

The pattern explains placement decisions. Tensor and expert parallel collectives sit on the critical path of every layer and move activation-sized data many times per step, so they are kept inside a node on the fastest links, NVLink and NVSwitch (see NVLink and NVSwitch architecture). Data-parallel all-reduce happens once per step and can be overlapped with backward compute, so it tolerates the slower inter-node network. Pipeline send and receive moves small amounts between neighbours and is the cheapest to stretch across nodes.

A cost model you can plan with

A simple model predicts collective time well enough to plan with. Each message costs a fixed latency alpha plus its size divided by bandwidth B. A ring all-reduce over n ranks takes 2(n-1) steps, each sending S/n bytes:

T_allreduce(ring)  ~= 2(n-1) * alpha  +  2(n-1)/n * S / B
T_allgather(ring)  ~=  (n-1) * alpha  +   (n-1)/n * S / B   # S = full gathered size
T_reducescatter    ~=  (n-1) * alpha  +   (n-1)/n * S / B   # S = full input size

For large tensors the bandwidth term dominates and the ring is close to optimal. For small tensors the latency term dominates, and it grows linearly with n. Tree and hierarchical algorithms reduce the number of latency steps to roughly logarithmic in n, which is why libraries switch to them for small messages and large rank counts. It is a simplification to call trees simply latency-optimal: well-pipelined tree algorithms can also reach high bandwidth, and the library picks per message size and topology. Treat the formulas as bounds, then measure.

Two practical consequences follow. First, many small collectives are expensive, which is why DDP groups gradients into buckets (25 MB by default) instead of reducing each parameter separately. Second, the bandwidth B in a hierarchical cluster is the slowest link the collective crosses: a data-parallel group that spans nodes runs at network speed, not NVLink speed.

The collectives in PyTorch code

PyTorch exposes each collective in torch.distributed. The calls below are what DDP, FSDP and MoE layers issue under the hood; writing them by hand once is the fastest way to understand a trace.

import os, torch, torch.distributed as dist

dist.init_process_group("nccl")                      # launched with torchrun
rank, n = dist.get_rank(), dist.get_world_size()
torch.cuda.set_device(int(os.environ["LOCAL_RANK"]))
S = 64 * 1024 * 1024                                 # elements per full tensor
x = torch.full((S,), float(rank), device="cuda", dtype=torch.bfloat16)

dist.all_reduce(x, op=dist.ReduceOp.SUM)             # DDP gradients

shard = torch.empty(S // n, device="cuda", dtype=torch.bfloat16)
dist.reduce_scatter_tensor(shard, x, op=dist.ReduceOp.SUM)   # FSDP gradients
full = torch.empty(S, device="cuda", dtype=torch.bfloat16)
dist.all_gather_into_tensor(full, shard)             # FSDP parameters

out = torch.empty_like(x)
dist.all_to_all_single(out, x)                       # MoE dispatch, equal splits

def busbw(op, nbytes, seconds):
    factor = {"all_reduce": 2 * (n - 1) / n, "all_gather": (n - 1) / n,
              "reduce_scatter": (n - 1) / n, "all_to_all": (n - 1) / n}[op]
    return nbytes / seconds * factor / 1e9           # GB/s

torch.cuda.synchronize(); t0 = torch.cuda.Event(enable_timing=True)
t1 = torch.cuda.Event(enable_timing=True)
for _ in range(5): dist.all_reduce(x)                # warm up
t0.record()
for _ in range(20): dist.all_reduce(x)
t1.record(); torch.cuda.synchronize()
sec = t0.elapsed_time(t1) / 1000 / 20
if rank == 0:
    print(f"all_reduce algbw {x.numel()*2/sec/1e9:.1f} GB/s  "
          f"busbw {busbw('all_reduce', x.numel()*2, sec):.1f} GB/s")

The last block uses the two bandwidth numbers from NVIDIA's nccl-tests. Algorithm bandwidth is simply size over time. Bus bandwidth multiplies it by the per-operation factor (2(n-1)/n for all-reduce, (n-1)/n for all-gather, reduce-scatter and all-to-all, 1 for broadcast and reduce), so the result is comparable to the hardware link speed and stays roughly constant as n changes. For all-gather and reduce-scatter, nccl-tests defines S as the full array size, not the shard.

Measuring your fabric with nccl-tests

Before trusting a model of your training step, measure the fabric. nccl-tests builds one binary per collective and sweeps message sizes:

# 8 GPUs in one process, sizes 8 B to 8 GB, doubling each step
./build/all_reduce_perf     -b 8 -e 8G -f 2 -g 8
./build/reduce_scatter_perf -b 8 -e 8G -f 2 -g 8
./build/alltoall_perf       -b 8 -e 1G -f 2 -g 8
# multi-node: one rank per GPU under MPI (build with MPI=1)
mpirun -np 16 -N 8 ./build/all_reduce_perf -b 1M -e 4G -f 2 -g 1

Read the output as a curve. At small sizes time is flat (latency bound); at large sizes bus bandwidth plateaus. The plateau is your B for planning; the size at which you reach about half of it tells you how big a bucket must be to be efficient. Run it on every new cluster, after driver or NCCL upgrades, and when a job is slower than its model predicts. A node whose plateau is well below its peers usually has a degraded link or a misconfigured network adapter.

Worked example: a 7B model on 16 GPUs

Consider a 7-billion-parameter model trained with BF16 gradients on 16 GPUs across two nodes. The gradient tensor is 7e9 x 2 bytes = 14 GB. Suppose nccl-tests across both nodes plateaus at a bus bandwidth of 40 GB/s, limited by the network.

DDP. One all-reduce of 14 GB per step. Each rank sends 2 x 15/16 x 14 GB = 26.25 GB, and the time is algorithm size over algorithm bandwidth, where algbw = busbw / factor = 40 / 1.875 = 21.3 GB/s, so 14 / 21.3 = about 0.66 s. If backward takes 1.5 s and the buckets overlap well, most of it is hidden; if not, it adds 0.66 s to every step.

FSDP (full sharding). Per step: all-gather parameters for forward (14 GB of BF16 parameters), all-gather again for backward, reduce-scatter gradients. Each moves 15/16 x 14 GB per rank, so the total is 3 x 13.1 = 39.4 GB, 1.5 times DDP's traffic. In exchange each GPU stores only 1/16 of parameters, gradients and optimizer state. The extra traffic is the price of memory, and it is why hybrid sharding (shard inside a node, replicate across nodes) is popular: the parameter gathers stay on NVLink and only gradient reduction crosses the network.

MoE layer. With 4,096 tokens of hidden size 4,096 in BF16 per rank, one all-to-all moves about 32 MB, twice per layer per pass. That is small, so latency and load imbalance dominate: if routing sends twice the average tokens to one expert, every rank waits for that rank.

Failure modes

  • Hangs from mismatched calls. Every rank must issue the same collectives in the same order with matching sizes and dtypes. A conditional all_reduce on one rank, or a parameter that only some ranks use, deadlocks the job. Set a timeout in init_process_group(timeout=...) so a hang becomes an error, and use PyTorch's NCCL flight recorder to see which collective each rank was in.
  • Silent slowness from small messages. Thousands of tiny collectives each pay latency. Bucket, fuse, or batch them.
  • Stragglers. A collective finishes when the slowest rank arrives. One thermally throttled GPU or a noisy neighbour slows everyone; the profile shows long waits inside the collective on the healthy ranks.
  • Wrong group placement. A tensor-parallel group that spans nodes runs its per-layer all-reduces over the network and can halve throughput. Check rank-to-node mapping.
  • All-to-all imbalance. Uneven expert routing makes all_to_all_single with uneven splits wait on the busiest rank. Monitor tokens per expert.
  • No overlap. Collectives serialized after compute. See Collective Communication Overlap for how to hide them.
  • Debugging blind. Run with NCCL_DEBUG=INFO once per new setup to confirm which transports and interfaces NCCL picked.

Trade-offs

Every choice trades memory, bandwidth and latency. DDP moves the least data but replicates everything. FSDP saves memory at 1.5 times the traffic and adds per-layer gathers that must be prefetched to hide. Tensor parallelism cuts per-GPU activation and weight memory but puts collectives on every layer's critical path, so it only pays off on fast intra-node links. Expert parallelism scales parameters cheaply but makes step time depend on routing balance. Lower-precision communication (for example reducing gradients in BF16 instead of FP32) halves bytes but can affect convergence; measure loss curves before adopting it.

What to do next

  1. Write down every collective your training step issues, with tensor size and group.
  2. Run nccl-tests for all-reduce, reduce-scatter, all-gather and all-to-all on your cluster and record the plateau bus bandwidth inside and across nodes.
  3. Predict per-step communication time with the cost model and compare it to a profiler trace.
  4. Keep tensor and expert parallel groups inside a node; put data-parallel groups across nodes.
  5. Size DDP or FSDP buckets above the half-bandwidth message size.
  6. Set collective timeouts and enable the NCCL flight recorder before your first long run.
  7. Re-run the benchmark after every driver, NCCL or network change.
Key takeaway: Each parallelism strategy is a pattern of collectives: DDP all-reduces gradients, FSDP all-gathers parameters and reduce-scatters gradients, tensor parallelism all-reduces activations and MoE all-to-alls tokens. Know the bytes each moves, measure bus bandwidth with nccl-tests, keep latency-critical groups on fast links and make every rank issue the same collectives in the same order.