PidokuInfra

Collectives and NCCL

Advanced 1h 15m Difficulty 4/5 Topic 07 of 12

Prerequisites 03, II.09


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        synchronization

2. 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 transparently

A 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:

  1. The bandwidth term approaches 2S/bandwidth as N grows — it does not grow with N. Ring AllReduce is bandwidth-optimal.
  2. 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 µs

LLM 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 of 2(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#

Shell
# 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 8

Typical 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.2

Read 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#

Shell
# 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 P2P

Set 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-tests on 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=INFO once 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=INFO will tell you.
  • Monitor for NCCL errors in logs. unhandled system error and remote process exited are 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#

  1. Explain ring AllReduce and derive its 2(N-1)/N data volume.
  2. Why is LLM decode latency-bound rather than bandwidth-bound in its collectives?
  3. What is the LL protocol and when does NCCL use it?
  4. What is a custom all-reduce and why is it faster for small messages?
  5. Why must you set NCCL_TIMEOUT in production?
  6. How would you verify that NVLink is actually being used?
  7. What happens when one rank in an NCCL group dies?

12. Further reading#

  • [REFERENCE] NCCL documentation and the nccl-tests repository
  • [FUNDAMENTAL] Baidu’s original ring-allreduce writeup
  • [REFERENCE] NCCL environment variable reference
  • [ESTABLISHED] vLLM and TensorRT-LLM custom all-reduce implementations
  • Next: 08 — Interconnects

↑↓ navigate↵ openesc close