← All writing
Reliability▶ Interactive

One dead link, 512 idle GPUs: ring all-reduce under failure

Ring all-reduce is close to bandwidth-optimal, and that same design makes it a chain. Break one link and every rank stalls, then every rank reports the same timeout.

Why a ring?

Data-parallel training ends every step by summing gradients across all GPUs. The naive approach sends everything to one GPU and broadcasts the result back, which makes that GPU’s link the bottleneck. Ring all-reduce arranges the N GPUs in a cycle and splits the buffer of size S into N chunks:

  1. Reduce-scatter (N−1 steps): each GPU sends one chunk to its right neighbour and adds the chunk arriving from its left. Afterwards, each GPU owns one fully-reduced chunk.
  2. All-gather (N−1 steps): the reduced chunks travel once more around the ring until every GPU has all of them.

Each step moves S/N bytes per link, and there are 2(N−1) steps. The standard α-β cost model gives:

T(N, S) ≈ 2(N−1)·α + 2·(N−1)/N · S / B

Here α is per-step latency and B is per-GPU link bandwidth. The bandwidth term approaches 2S/B as N grows, so it is independent of cluster size. The latency term grows linearly with N, which is why NCCL also has tree and other algorithms for small messages at scale. Try it:

All-reduce cost model α-β model, single ring
Total time
–
Bandwidth term
–
Latency term
–
busbw
–
nccl-tests definition

The price of being a chain

Nothing in the ring is optional. In step s, a GPU can only send what it finished in step s−1, and it can only finish step s−1 once its left neighbour’s data has arrived. Every rank depends on every other rank within one lap.

Run the ring below, then click any link to cut it mid-collective (or use the button). The stall spreads one hop per step until every rank is waiting. Then the watchdogs fire.

Ring all-reduce simulator Click a link to cut it
running blocked on recv watchdog timeout collective done

What every rank sees

Look at the log after the watchdogs fire. Every rank reports the same thing: a timeout in the same ALLREDUCE, with the same sequence number. The order is effectively random, because each rank’s watchdog thread wakes on its own schedule. The rank that reports first is just the one whose timer ran out first.

This is the core attribution problem in multi-node training. One bad cable produces N identical error messages, and the most visible one (the first in the aggregated logs) points at the wrong machine. Real clusters make it worse:

What attribution actually needs

To name the culprit, you need evidence that differs across ranks, not the error message that is identical on all of them:

The simulator’s step counters stand in for that flight recorder. Cut a link, let it run, and compare the step numbers: the lowest one sits just downstream of the cut. Try the incident drill to do the same with mixed signals.

Why it matters for recovery

Without attribution, the only safe move is to restart the whole job from the last checkpoint. Every rank pays for process start, CUDA and NCCL init, checkpoint load and warmup, plus the work since the last checkpoint, and the job may land on the same bad link again. With attribution, smaller recoveries become safe: drain one node, re-initialize communicators, and restore from a peer’s in-memory copy. How often you should checkpoint depends on all of this; there’s a calculator for that.

This is what I’m building in goodput: fault injection with ground truth, flight-recorder signature capture, cross-layer correlation, and topology-aware recovery. It’s at milestone 1 of 6, with the design and harness in public and no results yet.

Keep going

Next · interactive
Reading the crime scene: triaging an NCCL timeout
Repo
goodput: fault attribution & recovery ↗