Data-parallel training ends every step the same way: each GPU holds a gradient computed on its own slice of the batch, and every GPU needs the sum. Shipping all gradients to one server makes that server's link the bottleneck, growing with the number of workers. The ring all-reduce removes the bottleneck entirely: each GPU sends and receives about twice the gradient size no matter how many GPUs there are. Patarasuk and Yuan showed in 2009 that this traffic is optimal for all-reduce, Baidu brought the algorithm to deep learning in 2017, Horovod made it a library, and it remains one of the core algorithms inside NCCL.

This article builds the ring from the step schedule up: who sends which chunk when, an implementation you can run, the refinements real libraries add (pipelining, multiple rings, hierarchy), and how it fails in training jobs. For the alpha-beta cost model, trees, NCCL protocols and how to read busbw, see All-Reduce, in depth; this page stays inside the ring.

The schedule

Arrange p ranks in a logical ring, each sending only to its right neighbour and receiving only from its left. Split the N-element buffer on every rank into p chunks. The algorithm is two phases of p - 1 steps each.

Reduce-scatter. At step s, rank r sends chunk (r - s) mod p to rank r + 1, which adds it into its own copy of that chunk. Each chunk travels around the ring collecting one contribution per hop. After p - 1 steps, rank r holds the complete sum for chunk (r + 1) mod p.

All-gather. At step s, rank r sends chunk (r + 1 - s) mod p, and the receiver overwrites rather than adds. The finished chunks travel around the ring, and after another p - 1 steps every rank has every reduced chunk.

Every rank is busy sending and receiving at every step, every link carries one chunk of N/p elements per step, and there are 2(p - 1) steps, so each rank sends 2(p - 1)/p x N elements in total. That is under 2N for any p, which is why ring bandwidth cost is essentially flat as the job grows. The price is latency: 2(p - 1) sequential steps, each paying a fixed startup cost.

Ring all-reduce on 4 GPUs: reduce-scatter step 0GPU 0sends chunk 0GPU 1sends chunk 1GPU 2sends chunk 2GPU 3sends chunk 3each link carriesN/4 per stepPhase 1: reduce-scatter (3 steps)step s: GPU r sends chunk (r - s) mod 4receiver adds it into its own copyafter 3 steps GPU r owns the fullsum of chunk (r + 1) mod 4Phase 2: all-gather (3 steps)step s: GPU r sends chunk (r + 1 - s) mod 4receiver overwrites its copyPer GPU, sent and received:2 (p - 1) / p x N, independent of pfor large p; latency grows as 2 (p - 1)
The first step of reduce-scatter on four ranks. Every link is used in the same direction at once, which is why the ring saturates the slowest link and nothing else.

Tracing four GPUs

Make it concrete with p = 4 and N = 4, so each chunk is one number. Rank r starts with (r + 1) x [1, 10, 100, 1000]; the expected result on every rank is [10, 100, 1000, 10000]. The table lists the chunk index each rank sends at each step.

StepRank 0 sendsRank 1 sendsRank 2 sendsRank 3 sends
RS 0c0c1c2c3
RS 1c3c0c1c2
RS 2c2c3c0c1
AG 0c1 (final)c2 (final)c3 (final)c0 (final)
AG 1c0c1c2c3
AG 2c3c0c1c2

Follow chunk 0. At RS 0, rank 0 sends its 1 to rank 1, which now holds 1 + 2 = 3. At RS 1, rank 1 sends 3 to rank 2, which holds 3 + 3 = 6. At RS 2, rank 2 sends 6 to rank 3, which holds 6 + 4 = 10, the full sum. In the all-gather, rank 3 passes 10 to rank 0, rank 0 to rank 1, and rank 1 to rank 2. Each rank sent six elements, matching 2(p - 1)/p x N = 2 x 3/4 x 4.

An implementation

The simulator below reproduces the schedule on lists of arrays and is the right thing to unit-test first. Chunk boundaries handle N not divisible by p by giving the first N mod p chunks one extra element.

import numpy as np

def chunk_bounds(n, p):
    base, extra = divmod(n, p)
    starts = [i * base + min(i, extra) for i in range(p + 1)]
    return [(starts[i], starts[i + 1]) for i in range(p)]

def ring_allreduce_sim(buffers):
    p, n = len(buffers), len(buffers[0])
    bufs = [b.astype(np.float64).copy() for b in buffers]
    bounds = chunk_bounds(n, p)
    for s in range(p - 1):                          # reduce-scatter
        msgs = [((r + 1) % p, (r - s) % p, bufs[r][slice(*bounds[(r - s) % p])].copy())
                for r in range(p)]
        for dst, ch, data in msgs:                  # all sends happen "at once"
            bufs[dst][slice(*bounds[ch])] += data
    for s in range(p - 1):                          # all-gather
        msgs = [((r + 1) % p, (r + 1 - s) % p, bufs[r][slice(*bounds[(r + 1 - s) % p])].copy())
                for r in range(p)]
        for dst, ch, data in msgs:
            bufs[dst][slice(*bounds[ch])] = data
    return bufs

On real GPUs the same index arithmetic becomes paired point-to-point operations. The send and the receive of each step must be posted together; if every rank does a blocking send before its receive, every rank waits for a receiver that is itself waiting, and the ring deadlocks.

import torch
import torch.distributed as dist

def ring_allreduce_(t):
    """In-place sum over the default group. t must be contiguous."""
    p, r = dist.get_world_size(), dist.get_rank()
    if p == 1:
        return t
    chunks = list(torch.tensor_split(t.view(-1), p))    # views; first chunks largest
    right, left = (r + 1) % p, (r - 1) % p
    scratch = torch.empty_like(chunks[0])
    for s in range(p - 1):                              # reduce-scatter
        send_c, recv_c = (r - s) % p, (r - s - 1) % p
        buf = scratch[: chunks[recv_c].numel()]
        ops = [dist.P2POp(dist.isend, chunks[send_c], right),
               dist.P2POp(dist.irecv, buf, left)]
        for req in dist.batch_isend_irecv(ops):
            req.wait()
        chunks[recv_c].add_(buf)
    for s in range(p - 1):                              # all-gather
        send_c, recv_c = (r + 1 - s) % p, (r - s) % p
        ops = [dist.P2POp(dist.isend, chunks[send_c], right),
               dist.P2POp(dist.irecv, chunks[recv_c], left)]
        for req in dist.batch_isend_irecv(ops):
            req.wait()
    return t

This is a teaching implementation using the same index arithmetic as the simulator, and it is far slower than dist.all_reduce, which runs fused CUDA kernels that reduce while data streams in. Use it to check your understanding against the library, comparing results with torch.allclose rather than exact equality, since summation order differs.

What real libraries add

Production rings add four refinements to the textbook schedule.

Pipelining. Each step moves a whole chunk of N/p elements, and a rank cannot forward a chunk until it has arrived. Libraries split chunks into smaller slices so a rank forwards the first slice while the next is still arriving; the ring then behaves like a pipeline, and step latency overlaps transfer.

Bidirectional rings. Links are full duplex. Running a second ring in the opposite direction on half the data uses the reverse direction of every link and roughly halves the time when the link, not the GPU, is the limit.

Multiple rings, or channels. A GPU with several NVLinks or NICs cannot saturate them all with one ring. NCCL splits a collective across channels, each its own ring over different links, with roughly one thread block per channel. NCCL_DEBUG=INFO prints the ring order each channel uses, and NCCL_MIN_NCHANNELS and NCCL_MAX_NCHANNELS bound the count.

Hierarchy. Inside a node, links are fast; between nodes they are slower. A 2D scheme runs reduce-scatter inside each node, then an all-reduce across nodes on each GPU's 1/g shard among GPUs with the same local rank, then an all-gather inside the node. Each GPU's cross-node traffic shrinks by the factor g, the GPUs per node, and the cross-node rings run in parallel through separate NICs.

Worked example: flat versus hierarchical

Estimate one gradient all-reduce for a 7-billion-parameter model with bf16 gradients, so N is 14 GB. Assume, purely for the arithmetic, 100 GB/s per direction between GPUs in a node, 25 GB/s per GPU between nodes, and a per-step startup cost alpha of 10 to 20 microseconds. Substitute your measured numbers before relying on any of this.

SetupBandwidth termLatency termApproximate total
8 GPUs, one node, flat ring2 x 7/8 x 14 / 100 = 0.245 s14 steps x 10 us0.25 s
64 GPUs, 8 nodes, one flat ring2 x 63/64 x 14 / 25 = 1.10 s126 steps x 20 us1.10 s
64 GPUs, hierarchical0.1225 + 0.1225 + 0.1225 ssmall0.37 s

The flat ring is only as fast as its slowest hop, and with one ring crossing node boundaries every chunk eventually squeezes through a 25 GB/s link. In the hierarchical version the intra-node reduce-scatter costs 7/8 x 14 / 100, the cross-node all-reduce on a 1.75 GB shard among 8 nodes costs 2 x 7/8 x 1.75 / 25, and the intra-node all-gather costs the same as the reduce-scatter: 0.1225 s each. Real NCCL rings use many channels through many NICs and narrow the gap, which is why you measure with nccl-tests as described in GPU Collective Operations. The deeper lesson holds regardless: put the slow links in the smallest part of the algorithm. And in practice most of this time is hidden behind the backward pass by bucketed overlap, covered in overlapping communication with compute.

Failure modes

  • Deadlock from blocking point-to-point. Hand-written rings that send then receive hang immediately. Post both together.
  • One slow link or rank stalls everyone. Every step is synchronous around the ring, so a throttled GPU, a degraded NVLink or a busy NIC sets the speed for the job. Compare per-rank step times, not averages.
  • A dead rank hangs the job. Neighbours wait forever for a message. Set a timeout in init_process_group and let the watchdog fail fast so the job can restart from a checkpoint.
  • Ring order fights the topology. A ring that crosses sockets or PCIe switches more than necessary runs at the slowest crossing. Check the printed ring order.
  • Mismatched buffers. Ranks issuing all-reduce on tensors of different size, dtype or order hang or corrupt silently. Bucket order must match across ranks.
  • Small messages. With 2(p - 1) latency-bound steps, the ring is poor for small tensors at large p; libraries switch to trees there.
  • Overflow and non-determinism. Summing fp16 gradients from many ranks can overflow; pre-divide or accumulate in wider types. Summation order varies by chunk and world size, so results are not bitwise identical across different p.

Trade-offs

Ring all-reduce is bandwidth-optimal and topology-simple, and it dominates for large messages. Its latency grows linearly with p, which is why trees win for small messages and very large jobs, and why switch-based in-network reduction can beat both where available. Hierarchical rings trade one extra phase for far less traffic on slow links. Reduce-scatter and all-gather, the two halves of the ring, are also used separately: sharded optimisers and FSDP issue them independently, which is the communication pattern analysed in the DDP communication math.

You rarely choose the algorithm by hand. NCCL picks ring or tree per call from message size, rank count and topology, and NCCL_ALGO exists to override that choice when you are diagnosing a regression, not as a permanent tuning knob. The decision that is genuinely yours is upstream: how large your gradient buckets are, whether gradients are communicated in bf16 or fp32, and how ranks are placed on nodes. Larger buckets push each call into the ring's bandwidth-bound regime; placing ranks so that consecutive ring neighbours share a node keeps most hops on the fast links.

What to do next

  1. Run the simulator against a numpy sum for several p and N, including N not divisible by p.
  2. Trace the p = 4 table by hand once; it makes NCCL logs readable.
  3. Run the torch ring against dist.all_reduce on two to four GPUs and compare with torch.allclose.
  4. Run nccl-tests all_reduce_perf on your cluster within one node and across nodes, and record bus bandwidth by message size.
  5. Read the ring order from NCCL_DEBUG=INFO output and confirm it follows your physical topology.
  6. Plug your measured bandwidths into the worked-example table to predict step time, then check how much is hidden by overlap.
  7. Set a process-group timeout and test that a killed rank fails the job instead of hanging it.
Key takeaway: Ring all-reduce is a reduce-scatter followed by an all-gather, each p - 1 steps around a ring in which every rank sends one chunk of N/p per step, so per-GPU traffic stays under 2N regardless of scale. Its cost is latency linear in p and sensitivity to the slowest link, which real libraries address with pipelined slices, multiple channels and hierarchical rings.