NCCL and collective operations
Sum a 16 MiB gradient across eight GPUs the obvious way, by sending every buffer to GPU 0, adding them there and sending the answer back, and GPU 0 moves 224 MiB while each of its seven neighbours moves 32 MiB. Add GPUs and the root's share grows without limit. That picture is where "at scale you are communication bound" comes from, and the arithmetic behind it is right for the algorithm nobody uses.
A ring all-reduce has every GPU move the same amount, 2(N-1)/N times the buffer: 28 MiB at eight ranks, a shade under 32 MiB at sixty-four, and never 32 no matter how many you add. Nothing funnels. This page builds that collective with NCCL, checks it two ways, times it, and pulls it apart into the two simpler collectives it is made of.
This day needs two GPUs. Kaggle's T4 x2 is the only free two-GPU tier; the program ships a one-rank fallback for everyone else, and says what that fallback does not prove.
Five collectives, and the one that is two of the others
A collective is an operation every rank calls together, with the same
shape, and which is not finished until all of them have arrived. NCCL
gives you five, and their definitions are worth reading in NVIDIA's own
words rather than mine. Broadcast "copies an N-element buffer from the
root rank to all the ranks". Reduce combines a value from every rank and
"stores the result only in the receive buffer of a specified root rank".
AllReduce does the same combination but "stores the result in the receive
buffer of every rank", so that out[i] = in0[i]+in1[i]+…+in(k-1)[i] on
each of the k ranks. AllGather "gathers N values from k ranks into an
output buffer of size k*N, and distributes that result to all ranks".
ReduceScatter reduces like Reduce, "except that the result is scattered in
equal-sized blocks between ranks" (all five quoted from
https://docs.nvidia.com/deeplearning/nccl/user-guide/docs/usage/collectives.html
, checked 2026-09-01).
Read the last two together and the all-reduce falls out. Cut the buffer into N chunks. Reduce-scatter leaves rank r holding the finished sum of chunk r and nothing else. All-gather then hands every rank a copy of every finished chunk. Sum, then share: that is an all-reduce, assembled from two operations that each cost the same.
The cost is where the naive picture dies. In each phase a rank passes one chunk to its neighbour, N-1 times, so it sends (N-1)/N of the buffer per phase and 2(N-1)/N over both. Every rank sends that, in parallel, on its own outgoing link, which is why the number is per rank rather than divided among them. At N = 2 it is exactly one buffer each. At N = 8 it is seven quarters. It approaches two and stops.
The call that returns before it has done anything
Everything asynchronous so far in this course queued work and then got out
of the way. cudaMemcpyAsync puts a copy on a stream
and returns; the copy is ordered, and day 51 is
built on that being reliable. An NCCL collective is asynchronous in the
same sense, "returns when the operation has been effectively enqueued to
the given stream", and the work then runs on the device
(https://docs.nvidia.com/deeplearning/nccl/user-guide/docs/usage/streams.html
, checked 2026-09-01).
The habit that breaks is assuming the call gets you to that point on its own. It cannot, because a collective involves ranks this thread has not reached yet. The user guide is blunt: "every NCCL call may have to block, waiting for other threads/ranks to arrive, before effectively posting the NCCL operation on the given stream", so a plain loop over devices "could block on the first call waiting for the other ones" (https://docs.nvidia.com/deeplearning/nccl/user-guide/docs/usage/groups.html , checked 2026-09-01).
for (int r = 0; r < ranks; ++r) {
cudaSetDevice(r);
ncclAllReduce(d_grad[r], d_out[r], n, ncclFloat, ncclSum,
comms[r], streams[r]); // rank 0 waits here, forever
}
ncclGroupStart() and ncclGroupEnd() fix it by deferring: inside a
group the calls "can return without having enqueued the operation on the
stream", and everything is posted when the group closes. Which is also
why cudaStreamSynchronize "can therefore be called only after
ncclGroupEnd returns". The same rule holds one level up, with no group to
help you: if any rank skips a collective the others are all called, the
job stops with no error and no output. That is day
65's hang arriving in a place where the missing
participant is a whole GPU.
What the program has to prove before it prints a bandwidth
Full program in
code/day92-nccl/allreduce_grad.cu.
Every library call is checked in the same shape as every CUDA call. A
collective that returns ncclInvalidUsage and goes unread does not fail
there. It fails later, as a hang on an unrelated line, with the real cause
several hundred microseconds in the past.
#define NCCL_CHECK(call) \
do { \
ncclResult_t err_ = (call); \
if (err_ != ncclSuccess) { \
std::fprintf(stderr, "NCCL error %s:%d: %s: %s\n", __FILE__, \
__LINE__, #call, ncclGetErrorString(err_)); \
std::exit(EXIT_FAILURE); \
} \
} while (0)
One call builds the whole clique. ncclCommInitAll "creates a clique
of communicators (single process version) in a blocking way"
(https://docs.nvidia.com/deeplearning/nccl/user-guide/docs/api/comms.html
, checked 2026-09-01), which is the shape for one process driving several
devices. Each rank gets its own stream, its own gradient and its own
output buffer.
NCCL_CHECK(ncclCommInitAll(comms.data(), ranks, devs.data()));
for (int r = 0; r < ranks; ++r) {
const size_t k = static_cast<size_t>(r);
CUDA_CHECK(cudaSetDevice(r));
CUDA_CHECK(cudaStreamCreate(&streams[k]));
CUDA_CHECK(cudaMalloc(&d_grad[k], bytes));
CUDA_CHECK(cudaMalloc(&d_out[k], bytes));
fillLocalGradient<<<blocks, kThreadsPerBlock, 0, streams[k]>>>(
d_grad[k], kElems, r);
CUDA_CHECK(cudaGetLastError());
}
The synchronise sits below the group, never inside it.
NCCL_CHECK(ncclGroupStart());
for (int r = 0; r < ranks; ++r) {
const size_t k = static_cast<size_t>(r);
NCCL_CHECK(ncclAllReduce(d_grad[k], d_out[k], kElems, ncclFloat,
ncclSum, comms[k], streams[k]));
}
NCCL_CHECK(ncclGroupEnd());
for (int r = 0; r < ranks; ++r) {
CUDA_CHECK(cudaSetDevice(r));
CUDA_CHECK(cudaStreamSynchronize(streams[r]));
}
Two checks, because one of them would pass on a reduce. Each rank's
output is compared against a host sum in double of the buffers the
devices produced, at a tolerance derived from the reduction depth. Then
every rank's output is compared byte for byte against rank 0's. Drop the
second and a ncclReduce typo grades as correct on rank 0 and is never
looked at again.
The bandwidth line divides the work into an honest number and a
comparable one. nccl-tests defines the correction as
B = S/t * (2*(n-1)/n) = algbw * (2*(n-1)/n)
(https://github.com/NVIDIA/nccl-tests/blob/master/doc/PERFORMANCE.md ,
checked 2026-09-01), and the reason to care is that algorithm bandwidth
moves when you change the rank count even if the hardware is doing exactly
the same thing.
if (ranks == 1) {
std::printf("\none-rank API path: %.3f ms\n", slowestMs);
std::printf("algorithm bandwidth: n/a, no collective traffic\n");
std::printf("bus bandwidth: n/a, 2(N-1)/N is 0 at one rank\n");
} else {
const double seconds = static_cast<double>(slowestMs) * 1.0e-3;
const double algBw = static_cast<double>(bytes) / seconds / 1.0e9;
const double factor =
2.0 * static_cast<double>(ranks - 1) / static_cast<double>(ranks);
std::printf("\nslowest rank: %.3f ms, algorithm bandwidth %.1f GB/s\n",
slowestMs, algBw);
std::printf("bus bandwidth: %.1f GB/s (2(N-1)/N = %.3f)\n",
algBw * factor, factor);
}
The timing is per device, one event pair each, and the reported figure is the slowest rank's. There is no single stream that sees the whole collective, and a collective is not done until every rank is done, so an average across ranks would flatter it.
Results
Partially re-run, not verified as a multi-GPU lesson. The one-rank
fallback passed again on a Tesla T4 with driver 580.173.02 and CUDA 13.0
(V13.0.88), using the same pinned NCCL 2.30.7.1 stack, and exited 0. Its mean
API-path time was 0.142 ms, but no collective traffic exists at
one rank, so both bandwidth fields correctly report n/a. The required
Kaggle T4 x2 transcript is still blocked on access to two GPUs in one
machine. The full Day 92 claim is therefore INCONCLUSIVE.
| Ranks | Slowest rank (ms) | algbw (GB/s) | 2(N-1)/N | busbw (GB/s) |
|---|---|---|---|---|
| 1 (fallback) | 0.142 | n/a | 0 | n/a |
| 2 (Kaggle T4 x2) | pending | pending | 1 | pending |
Four things the first two-GPU run has to settle, each of which can lose:
- Held for the fallback only. The one rank agreed with the double-precision host sum over all 4,194,915 elements. The cross-rank byte comparison requires two ranks and remains pending. On that run, every rank must agree with the double-precision host sum, and the two output buffers must be byte-identical. A pass on the first check and a failure on the second means the collective reduced but did not distribute, and the program names which rank diverged.
- Pending. At two ranks, bus bandwidth should equal algorithm bandwidth, because 2(N-1)/N is exactly 1 there. If the two columns disagree the arithmetic is wrong, not the hardware, and that is the cheapest bug in this program to find.
- Pending. The two-rank measured rate should land below the
host-to-device
bandwidth day 53
measured on one card. Kaggle's two T4s have no
NVLink, so their traffic crosses
PCIe, and a peer transfer has a device on both ends
of a link that a pinned host copy uses in one direction. If it lands
above, the transport is not what this prediction assumes and
NCCL_DEBUG=INFOsays what it is instead. - VERIFIED for the fallback contract. It passed every applicable
correctness check and reported both algorithm and bus bandwidth as
n/a. With a single contribution the sum is the input, so this proves API wiring and the buffer contract only.
The number that carries to another machine is bus bandwidth, not milliseconds. That is what the correction factor exists for: it is the same figure at two ranks and at sixty-four when the links are the same.
Run it yourself
Kaggle, accelerator set to T4 x2, is the free path and the only one
(FACT-SHEET.md section 4). Never the free P100: it is sm_60, which
CUDA 13 cannot compile for at all. Kaggle's image already carries NCCL,
so the build is one nvcc line with -lnccl. There is no Compiler
Explorer embed for this day, and not for the usual size reason: the CE
image has no NCCL headers, so the source fails there at
fatal error: nccl.h: No such file or directory, captured in the
README. On one GPU, run it anyway for
the fallback path, and see free GPU tiers for
what else is on offer.
Exercise
Replace the ncclAllReduce call with an ncclReduceScatter followed by
an ncclAllGather over the same buffers, and report whether the two paths
agree bitwise or only within the tolerance.
Time: 30 to 45 minutes. Submit: the modified file, the two timings, and one sentence saying which kind of agreement you got and why.
Check: the harness runs your two-phase version against the same
double-precision host sum and the same cross-rank byte comparison as the
shipped program, then diffs your output against the all-reduce output
element by element. It prints the first differing index, both values, and
whether the difference is inside the tolerance. A chunk-size mistake
usually shows as every element after the first n/ranks being wrong,
which the index tells you at a glance.
Hint 1
The scatter phase leaves each rank holding a different piece of the answer, and the gather phase has to put those pieces back in the right order. What size is the piece, and what does each call's count argument mean: the whole buffer, or the piece?
Hint 2
ncclReduceScatter takes recvcount and ncclAllGather takes
sendcount, and both are the per-rank chunk, not the full buffer. Your
intermediate buffer only needs to hold one chunk. Also check what happens
when kElems is not divisible by the rank count.
Solution
Both calls take the chunk, kElems / ranks, and the intermediate buffer
is that size. The two phases go inside one group each, or inside one group
together, and the output must equal the all-reduce output.
On whether it is bitwise: expect it, but do not assume it. If NCCL picked its ring algorithm for both runs the additions happen in the same order and the results are identical bit for bit; if it picked a tree for the one-shot all-reduce, the order differs and only the tolerance saves you. Either answer is correct and the interesting part is knowing which you got.
A collective's cost is set by the algorithm underneath it, and a decomposition that costs the same is a decomposition you are allowed to substitute. That is the licence behind gradient bucketing, overlap and every other trick in distributed training.
Pitfalls
The job stops and prints nothing. Every rank must call the same collective, the same number of times, with the same count, datatype and operation. One rank taking a branch that skips an all-reduce leaves the others waiting for a participant that will never arrive, and NCCL has nothing to report because nothing failed. Day 65 covers finding a hang; here the first thing to check is whether the counts match on every rank.
One thread, several devices, no group calls. The loop hangs on the first collective. The group is not an optimisation and there is no version of this shape that works without it.
A cudaStreamSynchronize between ncclGroupStart and
ncclGroupEnd. Inside the group nothing has been enqueued yet, so the
sync returns immediately on an empty stream and the code after it reads a
buffer that has not been written.
The first collective is slow and you timed it. NCCL builds its channels and loads its kernels on the first call of a shape. The program runs three warm-up collectives for the same reason day 11 warms up every kernel it times.
Comparing algorithm bandwidth across different rank counts. It moves with N even when the hardware does not. Multiply by 2(N-1)/N and compare bus bandwidth, which is what the correction is for.
Expecting the same bits from a different rank count. The reduction order changes with the topology and the algorithm, so a sum over four ranks need not equal the same sum over two in the last bit. Day 68 covers what reproducibility you can actually ask for.
Go deeper
- NCCL user guide, "Collective Operations": https://docs.nvidia.com/deeplearning/nccl/user-guide/docs/usage/collectives.html (checked 2026-09-01)
- NCCL user guide, "Group Calls", including management of multiple GPUs from one thread: https://docs.nvidia.com/deeplearning/nccl/user-guide/docs/usage/groups.html (checked 2026-09-01)
- NCCL API, communicator creation: https://docs.nvidia.com/deeplearning/nccl/user-guide/docs/api/comms.html (checked 2026-09-01)
- nccl-tests, "Performance", where algbw and busbw are defined: https://github.com/NVIDIA/nccl-tests/blob/master/doc/PERFORMANCE.md (checked 2026-09-01)
Next
Day 91 split a vector add by hand and moved the
halves yourself; this page handed that job to a library that knows the
topology. Day 93 goes back to one GPU and revisits unified memory with
cudaMemAdvise and prefetch, which is the same question as this one
asked about the host link instead of the peer link: who moves the bytes,
and when.