Day 92Module 10
draft

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:

  1. 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.
  2. 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.
  3. 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=INFO says what it is instead.
  4. 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

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.