1. What is it?#
Collectives are communication operations involving all ranks in a group. NCCL (NVIDIA Collective Communications Library) is the implementation everything uses.
The operations you need to know:
BROADCAST one rank's data → all ranks
REDUCE all ranks' data → summed on one rank
ALLREDUCE all ranks' data → summed, result on ALL ranks ← TP uses this
ALLGATHER each rank's slice → concatenation on all ranks
REDUCESCATTER sum across ranks, then each rank keeps one slice
ALL-TO-ALL each rank sends distinct data to every other rank ← MoE uses this
SEND/RECV point-to-point ← PP uses this
BARRIER synchronization2. Why NCCL specifically#
Because naive implementations are catastrophically slow. NCCL:
✓ Uses topology awareness (NVLink vs PCIe vs network) to pick algorithms
✓ Implements ring and tree algorithms that achieve near-peak bandwidth
✓ Runs on the GPU (kernels), not the CPU — no host involvement
✓ Overlaps with compute (separate streams)
✓ Handles GPUDirect P2P and RDMA transparentlyA hand-rolled AllReduce using cudaMemcpy and host coordination is typically 5-20x slower.
3. The ring AllReduce — how it achieves optimal bandwidth#
This is worth understanding because it explains the 2(N-1)/N factor that appears in every
communication cost calculation.
N ranks, each with a buffer of S bytes, divided into N chunks.
PHASE 1 — REDUCE-SCATTER (N-1 steps)
Step k: rank i sends chunk (i-k) mod N to rank (i+1) mod N,
and adds the chunk it receives.
After N-1 steps: rank i holds the fully-reduced chunk i.
PHASE 2 — ALL-GATHER (N-1 steps)
Step k: rank i sends its complete chunk to rank (i+1) mod N.
After N-1 steps: every rank has every reduced chunk.
Data sent per rank: 2 × (N-1)/N × S
Time: 2 × (N-1)/N × S / bandwidth + 2 × (N-1) × latency
└────── bandwidth term ──────┘ └── latency term ──┘Two crucial observations:
- The bandwidth term approaches
2S/bandwidthas N grows — it does not grow with N. Ring AllReduce is bandwidth-optimal. - The latency term is
2(N-1) × αand does grow with N. At N=8 that’s 14 hops.
For SMALL messages, the latency term dominates:
N=8, S=64 KB, NVLink (α=5 µs, β=450 GB/s):
bandwidth term: 2 × (7/8) × 64 KB / 450 GB/s = 0.25 µs
latency term: 2 × 7 × 5 µs = 70 µs ← 280x larger!
For LARGE messages, the bandwidth term dominates:
N=8, S=64 MB:
bandwidth term: 2 × (7/8) × 64 MB / 450 GB/s = 249 µs
latency term: 70 µsLLM decode has small messages (batch × d × 2 bytes = a few hundred KB), so it is latency-bound in its collectives. This is why:
- NCCL uses tree algorithms for small messages (fewer hops:
2·log₂(N)instead of2(N-1)). - Engines ship custom one-shot all-reduce kernels for very small messages.
- TP scaling degrades past 8 — the latency term grows while the compute shrinks.
4. Tiny example — measuring it#
# Build nccl-tests
git clone https://github.com/NVIDIA/nccl-tests && cd nccl-tests && make
# AllReduce across 8 GPUs, message sizes from 1 KB to 1 GB
./build/all_reduce_perf -b 1K -e 1G -f 2 -g 8
# All-to-All (for MoE)
./build/alltoall_perf -b 1K -e 1G -f 2 -g 8Typical output on an 8×H100 NVLink node:
size count type time(us) algbw(GB/s) busbw(GB/s)
1K 256 float 8.9 0.11 0.20
16K 4096 float 9.4 1.74 3.05
256K 65536 float 14.2 18.5 32.3
4M 1048576 float 78.3 53.5 93.6
64M 16777216 float 881.0 76.2 133.3
1G 268435456 float 13120.0 81.8 143.2Read the time(us) column for small sizes: 8.9 µs for 1 KB, 9.4 µs for 16 KB. The time is
essentially constant — it’s pure latency. That flat region is where LLM decode lives.
busbw (bus bandwidth) accounts for the 2(N-1)/N factor and is the number to compare against
hardware peak. 143 GB/s busbw on NVLink 4… that’s per-GPU and reflects the ring structure.
Run this on your hardware and record the results. You will use them in every communication cost estimate.
5. The environment variables that matter#
# DEBUGGING — start here when things are wrong
NCCL_DEBUG=INFO # prints topology detection and algorithm choice
NCCL_DEBUG_SUBSYS=INIT,GRAPH # more detail on setup
# CORRECTNESS / RELIABILITY
NCCL_TIMEOUT=1800 # seconds before a hung collective aborts
# DEFAULT IS VERY LONG — set this!
TORCH_NCCL_BLOCKING_WAIT=1 # fail fast instead of hanging
TORCH_NCCL_ASYNC_ERROR_HANDLING=1
# NETWORK SELECTION (multi-node)
NCCL_SOCKET_IFNAME=eth0 # which interface for bootstrap
NCCL_IB_HCA=mlx5_0,mlx5_1 # which InfiniBand adapters
NCCL_IB_DISABLE=0 # 1 to force TCP (debugging only — slow)
NCCL_NET_GDR_LEVEL=PHB # GPUDirect RDMA aggressiveness
# ALGORITHM TUNING
NCCL_ALGO=Ring,Tree # restrict algorithms
NCCL_PROTO=Simple,LL,LL128 # LL/LL128 are low-latency protocols for
# small messages — important for decode
NCCL_MIN_NCHANNELS=4 # more channels = more parallelism
NCCL_MAX_NCHANNELS=32
# P2P
NCCL_P2P_DISABLE=0 # 1 disables NVLink P2P (debugging)
NCCL_P2P_LEVEL=NVL # require NVLink for P2PSet NCCL_TIMEOUT in production. The default allows a hung collective to hang for a very
long time, holding GPUs. A finite timeout converts a hang into a crash, which your supervisor
can restart.
NCCL_DEBUG=INFO on first deployment to a new node type. It prints the detected topology
and chosen algorithms; if it says “using PCIe” where you expected NVLink, you’ve found a
problem before it becomes a mystery.
6. Under the hood — how NCCL picks an algorithm#
1. TOPOLOGY DETECTION at init: probes NVLink, PCIe, NUMA, network.
2. Builds a "graph" of possible rings/trees through the topology.
3. Per collective call, chooses:
- algorithm: Ring (bandwidth-optimal) vs Tree (latency-optimal)
- protocol: Simple / LL (low-latency, 8-byte flits) / LL128
- number of channels (parallel rings)
based on message size and the tuning model.
Rough thresholds (they vary by version and topology):
< 64 KB → Tree + LL (latency-optimized)
64 KB - 1 MB → Tree or Ring + LL128
> 1 MB → Ring + Simple (bandwidth-optimized)LL (“low latency”) protocol trades bandwidth for latency by using small flits with inline flags instead of separate synchronization — the right choice for LLM decode’s small messages.
If NCCL_DEBUG=INFO shows Simple protocol for your small decode AllReduces, something is
misconfigured.
7. Custom all-reduce for LLM inference#
Both vLLM and TensorRT-LLM ship custom all-reduce implementations for small messages:
NCCL ring AllReduce, N=8, 512 KB: ~35-50 µs
Custom one-shot AllReduce: ~12-20 µs
Mechanism: instead of a ring (14 hops), every GPU writes its data directly
into every other GPU's memory over NVLink (P2P), then each reads and reduces
locally. One "hop" instead of 2(N-1).
Only works when all GPUs are P2P-accessible (NVLink) and the message is small
enough that the N× redundant transfer is cheaper than the ring's hops.At 160 AllReduces per token, saving 25 µs each is 4 ms per token — very significant at small batch.
Verify your engine uses it. vLLM logs it; look for “custom allreduce” in the startup output. It’s disabled in some configurations (e.g. when P2P isn’t available, or above a size threshold).
8. Production implications#
- Measure your collectives with
nccl-testson every new node type. Record latency and busbw curves. They’re inputs to every capacity estimate. - Set
NCCL_TIMEOUT. Hangs become crashes; crashes are recoverable. - Run
NCCL_DEBUG=INFOonce per node type to verify topology detection. - Enable custom all-reduce where supported.
- For multi-node: verify GPUDirect RDMA is active. Without it, every cross-node byte bounces
through host memory.
NCCL_DEBUG=INFOwill tell you. - Monitor for NCCL errors in logs.
unhandled system errorandremote process exitedare the common ones and usually mean a rank died. - NCCL has no fault tolerance. One dead rank kills the group. Supervise accordingly.
9. Common mistakes#
Not setting NCCL_TIMEOUT. Hangs hold GPUs indefinitely.
Assuming NVLink is being used. Check with NCCL_DEBUG=INFO. A misconfigured container or a
missing device can silently fall back to PCIe.
Optimizing bandwidth when you’re latency-bound. LLM decode’s collectives are small; more bandwidth doesn’t help, fewer hops does.
Not using the custom all-reduce. 4 ms per token left on the table at small batch.
Multi-node without GPUDirect RDMA. Roughly halves effective bandwidth and adds latency.
Wrong NCCL_SOCKET_IFNAME. NCCL bootstraps over the wrong interface and either fails or
runs slowly.
Expecting NCCL to handle a failed rank. It won’t.
10. Hands-on exercise#
A. Benchmark your collectives. Run all_reduce_perf for sizes 1 KB to 1 GB on your
hardware. Plot time vs size on log-log axes. Identify the latency-bound and bandwidth-bound
regions. Find the crossover. Record in numbers.md.
B. Compute the TP overhead. Using your measured AllReduce times, compute the communication overhead for a model you serve at TP=2, 4, 8, at batch 1 and batch 64. Compare to the measured difference between TP degrees.
C. Verify the topology. Run a TP job with NCCL_DEBUG=INFO. Read the output: which
transport is used between each rank pair? Does it match nvidia-smi topo -m?
D. Custom all-reduce. Measure decode ITL with and without your engine’s custom all-reduce (usually a flag). Quantify the difference at batch 1 and batch 64.
E. Protocol effect. Run all_reduce_perf with NCCL_PROTO=Simple and with the default.
Compare small-message latency. How much does LL buy?
F. Break it. Set NCCL_P2P_DISABLE=1 and re-measure. This simulates a PCIe-only topology.
Quantify what NVLink is worth for your workload.
11. Interview questions#
- Explain ring AllReduce and derive its
2(N-1)/Ndata volume. - Why is LLM decode latency-bound rather than bandwidth-bound in its collectives?
- What is the LL protocol and when does NCCL use it?
- What is a custom all-reduce and why is it faster for small messages?
- Why must you set
NCCL_TIMEOUTin production? - How would you verify that NVLink is actually being used?
- What happens when one rank in an NCCL group dies?
12. Further reading#
- [REFERENCE] NCCL documentation and the
nccl-testsrepository - [FUNDAMENTAL] Baidu’s original ring-allreduce writeup
- [REFERENCE] NCCL environment variable reference
- [ESTABLISHED] vLLM and TensorRT-LLM custom all-reduce implementations
- Next: 08 — Interconnects