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.
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.
| Step | Rank 0 sends | Rank 1 sends | Rank 2 sends | Rank 3 sends |
|---|---|---|---|---|
| RS 0 | c0 | c1 | c2 | c3 |
| RS 1 | c3 | c0 | c1 | c2 |
| RS 2 | c2 | c3 | c0 | c1 |
| AG 0 | c1 (final) | c2 (final) | c3 (final) | c0 (final) |
| AG 1 | c0 | c1 | c2 | c3 |
| AG 2 | c3 | c0 | c1 | c2 |
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 bufsOn 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 tThis 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.
| Setup | Bandwidth term | Latency term | Approximate total |
|---|---|---|---|
| 8 GPUs, one node, flat ring | 2 x 7/8 x 14 / 100 = 0.245 s | 14 steps x 10 us | 0.25 s |
| 64 GPUs, 8 nodes, one flat ring | 2 x 63/64 x 14 / 25 = 1.10 s | 126 steps x 20 us | 1.10 s |
| 64 GPUs, hierarchical | 0.1225 + 0.1225 + 0.1225 s | small | 0.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_groupand 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
- Run the simulator against a numpy sum for several p and N, including N not divisible by p.
- Trace the p = 4 table by hand once; it makes NCCL logs readable.
- Run the torch ring against
dist.all_reduceon two to four GPUs and compare withtorch.allclose. - Run nccl-tests all_reduce_perf on your cluster within one node and across nodes, and record bus bandwidth by message size.
- Read the ring order from
NCCL_DEBUG=INFOoutput and confirm it follows your physical topology. - Plug your measured bandwidths into the worked-example table to predict step time, then check how much is hidden by overlap.
- Set a process-group timeout and test that a killed rank fails the job instead of hanging it.