All-reduce takes one buffer on every participating GPU, combines them element-wise (usually a sum) and leaves the identical result on every GPU. It is the collective that makes data-parallel training work: each GPU computes gradients on its own slice of the batch, all-reduce sums them, and every replica applies the same update and stays in lockstep. In large jobs it moves more bytes than any other operation, so its cost often decides how far a model scales.

This article derives all-reduce from first principles: a cost model, the ring algorithm and why its bandwidth cost barely grows with GPU count, trees for latency, in-switch reduction, NCCL's protocols, how to read benchmark numbers, the APIs in CUDA C and PyTorch, a worked time estimate, numerics, failure modes and a checklist. Where this article says NCCL, the same ideas apply to AMD's RCCL, which follows the NCCL API.

A cost model for collectives

A simple model is enough to reason about collectives. Sending a message of n bytes over a link costs alpha + n / B, where alpha is the fixed latency per message (kernel launch, synchronisation, network hop, a few microseconds) and B is the link bandwidth. A collective algorithm is a schedule of such messages, and its cost is a latency term (alpha times the number of sequential steps) plus a bandwidth term (bytes on the busiest link divided by B). Reduction arithmetic is usually cheap on a GPU next to the transfers.

The naive approach shows why the schedule matters. Gather everything at GPU 0, sum, and broadcast back: GPU 0 receives (N - 1) S bytes and sends (N - 1) S bytes, so its links carry a load that grows linearly with N while every other GPU's links sit mostly idle. At 64 GPUs the root moves 63 copies of the gradient each way. Any good algorithm spreads the traffic so every link carries roughly the same load.

The ring algorithm

The ring algorithm arranges N GPUs in a cycle and splits each buffer into N chunks. It runs in two phases. In reduce-scatter, at step k every GPU sends one chunk to its right neighbour and adds the chunk arriving from its left neighbour into its own copy. Chunks rotate so that after N - 1 steps each GPU holds one chunk that contains the sum from all N GPUs. In all-gather, the GPUs pass those finished chunks around the ring for another N - 1 steps, overwriting rather than adding, until every GPU has all N summed chunks.

Count the traffic. Each step sends S / N bytes per GPU, and there are 2 (N - 1) steps, so each GPU sends 2 (N - 1) S / N bytes and receives the same. As N grows that tends to 2S, independent of N, and no algorithm can do much better: every GPU must receive at least (N - 1) S / N bytes of other GPUs' data to contribute to the result, and must receive the final values of the chunks it did not reduce, and a standard lower-bound argument turns that into the same 2 (N - 1) S / N figure. The ring is bandwidth-optimal. Its weakness is latency: 2 (N - 1) sequential steps, so for small messages on many GPUs the alpha term dominates. NCCL also runs several rings at once over different channels, each handling a slice of the buffer, to use all available links in parallel.

Ring all-reduce on 4 GPUs: reduce-scatter, then all-gatherGPU 0c0c1c2c3c0c1c2c3GPU 1c0c1c2c3c0c1c2c3GPU 2c0c1c2c3c0c1c2c3GPU 3c0c1c2c3c0c1c2c3After N-1 = 3 reduce-scatter steps each GPU owns one fully summed chunk (yellow).Each step sends one chunk of size S/N to the right neighbour and adds what arrives from the left.After N-1 = 3 all-gather steps every GPU holds every summed chunk (green).Per-GPU traffic: 2 (N-1)/N x S, almost independent of Nlatency grows with 2 (N-1) steps; bandwidth term stays near 2S / link bandwidth
The ring splits the buffer into N chunks; every link carries the same load and per-GPU traffic is 2 (N-1)/N of the buffer.
# Ring all-reduce simulated on lists; rank r sends to (r+1) % N.
def ring_all_reduce(bufs):
    N = len(bufs)
    chunks = [list(split(b, N)) for b in bufs]        # chunks[r][k]
    # reduce-scatter: after step t, rank r has accumulated chunk (r - t - 1) % N
    for t in range(N - 1):
        sends = [(r, (r - t) % N, chunks[r][(r - t) % N]) for r in range(N)]
        for r, k, data in sends:
            dst = (r + 1) % N
            chunks[dst][k] = [a + b for a, b in zip(chunks[dst][k], data)]
    # rank r now owns the full sum of chunk (r + 1) % N
    for t in range(N - 1):
        sends = [(r, (r + 1 - t) % N, chunks[r][(r + 1 - t) % N]) for r in range(N)]
        for r, k, data in sends:
            chunks[(r + 1) % N][k] = list(data)        # copy, no add
    return [sum(c, []) for c in chunks]

def split(b, N):
    n = (len(b) + N - 1) // N
    return (b[i * n:(i + 1) * n] for i in range(N))

Trees and in-network reduction

Trees attack the latency term. Reduce up a binary tree to the root and broadcast back down, and the number of sequential steps is about 2 log2 N instead of 2 (N - 1). With pipelining, the buffer is cut into many small pieces that flow up and down the tree concurrently, so the bandwidth cost stays reasonable too. A plain tree wastes half the links, because leaves only send up and receive down. NCCL's double binary tree, introduced in version 2.4, uses two complementary trees, each carrying half the data, arranged so that a node which is a leaf in one tree is an interior node in the other. Together they use every node's bandwidth and keep logarithmic latency, which is why trees tend to win for small and mid-sized messages across many nodes, while rings win for large messages.

Two further families move the reduction into the network. On systems with NVLink switches, NCCL's NVLS algorithms use the switch's in-network reduction (NVLink SHARP) so GPUs send their data once and the switch returns the sum, cutting the traffic each GPU must move. On InfiniBand fabrics with SHARP-capable switches, the CollNet algorithms do the same across nodes. Availability depends on hardware generation, firmware and NCCL version, so check what your NCCL actually selected rather than assuming.

Protocols: Simple, LL and LL128

Independently of the algorithm, NCCL chooses a protocol for how data and readiness flags travel. Simple moves large chunks and synchronises with memory fences, which gives the best bandwidth with the highest per-step latency. LL (low latency) packs 4 bytes of data with a 4-byte flag into every 8-byte store, so a receiver can poll the data itself without a fence; latency is low but half the bandwidth goes to flags. LL128 does the same on 128-byte lines with 120 bytes of data, keeping most of the bandwidth at low latency on hardware where 128-byte stores are delivered atomically. NCCL picks an algorithm and protocol per call from a tuning model based on message size, topology and GPU count. You can override it with NCCL_ALGO and NCCL_PROTO for experiments, but forcing them in production usually makes some message sizes slower.

Reading benchmarks: algbw and busbw

The standard benchmark is all_reduce_perf from the nccl-tests repository. It reports two bandwidths, and confusing them causes most misreadings. Algorithm bandwidth (algbw) is the buffer size divided by the time. Bus bandwidth (busbw) multiplies algbw by 2 (N - 1) / N for all-reduce, the ring's traffic factor, so it estimates the bandwidth each GPU's links actually sustained and can be compared with the hardware link speed regardless of N. If busbw at large sizes is far below what your links should deliver, something in the topology or configuration is wrong; if algbw falls as you add GPUs while busbw holds, that is just the 2 (N - 1) / N factor.

# Single node, 8 GPUs, sizes 8 B to 8 GB doubling each step
./build/all_reduce_perf -b 8 -e 8G -f 2 -g 8

# Multi-node with MPI: one process per GPU
mpirun -np 16 -N 8 -x NCCL_DEBUG=INFO ./build/all_reduce_perf -b 8 -e 8G -f 2 -g 1

The APIs

In CUDA C the call is asynchronous on a stream: it enqueues work and returns. Every rank must call the same collectives in the same order with the same count and type, or the job hangs. ncclAvg (available since NCCL 2.10) divides by N inside the collective, saving a separate scaling kernel.

// One rank per process; comm created with ncclCommInitRank.
ncclResult_t r = ncclAllReduce(grads, grads, count,   // in-place is allowed
                               ncclBfloat16, ncclAvg, comm, stream);
if (r != ncclSuccess) handle(ncclGetErrorString(r));
cudaStreamSynchronize(stream);   // only when the host needs the result

# PyTorch: the same collective through torch.distributed
import datetime
import torch.distributed as dist
dist.init_process_group("nccl", timeout=datetime.timedelta(minutes=10))
work = dist.all_reduce(grad_bucket, op=dist.ReduceOp.SUM, async_op=True)
# ... keep computing other gradients ...
work.wait()
grad_bucket /= dist.get_world_size()

You rarely call it by hand in training. DistributedDataParallel groups gradients into buckets (25 MB by default) and launches an all-reduce for each bucket as soon as backward produces it, overlapping communication with the rest of backward. FSDP and ZeRO replace the single all-reduce with its two halves: reduce-scatter for gradients, so each rank keeps only its shard, and all-gather for parameters before use. The ring decomposition above is exactly why that split costs no extra bandwidth.

Worked example: how long does a step's all-reduce take?

Estimate the cost for a 1-billion-parameter model trained with plain data parallelism and bf16 gradients: S = 2 GB per step. On one node of 8 GPUs, suppose nccl-tests shows a large-message busbw of 350 GB/s (use your own measurement; it depends on the GPU and NVLink generation). The time is about 2 (N - 1) / N x S / busbw = 1.75 x 2 / 350 = 10 ms. If the backward pass takes 150 ms and DDP overlaps most buckets, communication is nearly hidden.

Now go to 4 nodes, 32 GPUs, with each node attached to the fabric at an aggregate 400 GB/s but a measured inter-node busbw of, say, 80 GB/s for this job. Inter-node links are now the bottleneck: roughly 2 x 31 / 32 x 2 / 80, about 48 ms. That may no longer hide behind a 150 ms backward, especially for the last buckets, which only become ready at the end. The options are the ones the cost model suggests: increase compute per step (larger local batch, gradient accumulation), reduce bytes (sharded optimizers do not reduce all-reduce bytes, but lower-precision or compressed gradients do, at an accuracy risk), or improve the fabric configuration so busbw rises. NCCL already runs hierarchically, using NVLink inside the node and the network between nodes, so the first check is whether it found all NICs.

Numerics and reproducibility

Summation order affects floating-point results. Ring and tree sum the same numbers in different orders, so switching algorithms, or NCCL choosing differently because the GPU count changed, changes results in the last bits. Training is normally tolerant of this, but bitwise reproducibility across configurations is not something all-reduce provides. Summing bf16 gradients across many ranks loses precision, which is why mixed-precision setups often reduce gradients in fp32 or use a higher-precision accumulation where available, and why averaging (dividing by N) before or during the sum avoids overflow in fp16.

Failure modes

SymptomLikely causeWhat to do
Job hangs in all-reduce, no errorRanks called collectives in different orders, or one rank skipped a stepSet a process-group timeout; use PyTorch's flight recorder (TORCH_NCCL_TRACE_BUFFER_SIZE) to see each rank's last collective
busbw far below link speedWrong network interface, missing GPUDirect RDMA, PCIe crossing CPU socketsRun with NCCL_DEBUG=INFO and read the chosen transports; set NCCL_SOCKET_IFNAME / NCCL_IB_HCA
One slow rank slows everyoneThermal throttling, a degraded link, a noisy neighbourPer-rank timing; nccl-tests on node pairs to isolate the bad host
Results differ between runs at different scalesSummation order changed with the algorithmAccept, or pin NCCL_ALGO for reproducibility tests only
Small-message all-reduce dominatesLatency-bound: many tiny bucketsBigger buckets, fused gradients, tree algorithms
Errors after a NIC or switch faultNCCL communicators do not recover in placeRestart from checkpoint; use an elastic launcher

Trade-offs

Ring maximises bandwidth but pays linear latency; trees pay logarithmic latency with somewhat less bandwidth efficiency; in-network reduction needs specific switches but cuts both. Larger buckets amortise latency but delay the first all-reduce and reduce overlap. Gradient compression cuts bytes but adds compute and can hurt convergence. For most teams the practical order is: measure busbw, make sure NCCL sees the right topology, tune bucket size and overlap, and only then consider changing the parallelism strategy.

Related reading: the collective operations vocabulary, overlapping collectives with backward, NCCL collectives architecture, NVLink and FSDP.

What to do next

  1. Run all_reduce_perf on one node and across two nodes; record busbw at 1 GB and the size at which busbw saturates.
  2. Compare the measured busbw with your link specifications; if it is far off, read the NCCL_DEBUG=INFO output for transports, rings and NIC selection.
  3. Estimate your per-step all-reduce time with 2 (N - 1) / N x S / busbw and compare it with the backward pass duration.
  4. Profile a training step to see whether DDP buckets overlap backward; tune bucket size if the last buckets sit exposed.
  5. Set a process-group timeout and enable the flight recorder so a hang produces evidence instead of a stuck job.
  6. Decide on gradient reduction precision explicitly, and document that results are not bitwise reproducible across GPU counts.
Key takeaway: All-reduce leaves the element-wise sum of every GPU's buffer on every GPU, and the ring version does it with per-GPU traffic of 2 (N - 1) / N times the buffer, nearly independent of scale, while trees and in-network reduction trade some of that for lower latency. Measure busbw with nccl-tests, estimate step cost from it, keep collectives overlapped with backward, and add timeouts and tracing so a mismatched or stalled rank is diagnosable.