Skip to content

CMU 11-868 L14-L15: Distributed Training and Data Parallelism, and Where Gradient Sync Costs Come From

Sep 30, 20261 min
TL;DRLectures 14 and 15 of 11-868 go from the parameter server to PyTorch DDP. They use NCCL's five collectives (Broadcast, Reduce, AllReduce, ReduceScatter, AllGather) as building blocks, show why a ring makes broadcast time nearly independent of GPU count, and split AllReduce into ReduceScatter plus AllGather. The second lecture takes apart DDP's two key designs: bucketing gradients (25 MB by default) and starting synchronization before the backward pass finishes. There is no recording; this guide works from slide page numbers and the VLDB 2020 paper.

🌏 中文版

Version note: This post is based on the Spring 2026 edition of CMU 11-868 LLM Systems. Every fact was checked on 2026-09-30 against the course syllabus, the L14 slides (48 pages), the L15 slides (26 pages), and the PyTorch DDP paper. Access grade A3, but the course has no public recordings. Page numbers below always mean PDF pages; the numbers printed in the corner of the L14 slides run 2-3 higher than the PDF page, so go by the PDF.

Series: Previous: HW4: Fused CUDA Kernels for Softmax and LayerNorm | Next: L16-L17 Model Parallelism and MoE | Series overview

What these lectures answer

LightSeq in the previous lecture squeezed everything out of one GPU. The next question: when one isn't enough, where does the time go when many GPUs train together?

L14 slide 5 sets the scale. DeepSeek-V3 (671B) was pretrained on 2,048 H800s for about two months, 2.664 million H800 GPU hours in total; LLaMA 3.1 (405B) used 16,000 H100s and 30.84 million GPU hours. At that scale, how GPUs exchange data decides whether the money is well spent.

Slide 6 splits scaling strategies into two families:

Partition the dataPartition the model
Single-node data parallel, distributed data parallel, parameter serverModel parallel: pipeline parallel, tensor parallel

These two lectures cover only the left column: every GPU holds a full copy of the model and eats different data. The right column is the next post. The syllabus puts L14 on March 9 (the first class after spring break) and L15 on March 11, with the PyTorch DDP paper (VLDB 2020) as L15's reading. There is also a "Recitation 6 Distributed Training" on March 20, with no slides linked in the syllabus.

Intuition: data parallelism has only one thing to synchronize

Every data-parallel step looks the same (L14 slide 32):

  1. Split a batch across workers
  2. Each worker runs forward and backward on its share to get local gradients
  3. All workers compute the average gradient together
  4. Each worker updates its own copy of the parameters with the average

Steps 1, 2, and 4 need no communication at all. The only communication cost is step 3. So the whole problem becomes: how do you make that step fast enough that GPUs don't sit idle waiting?

Mechanism 1: from parameter server to AllReduce

The old way: parameter server

L14 slide 7 shows the classic parameter server. Workers pull parameters from the server, compute local gradients, push them back, and the server aggregates, updates, and starts the next round. Slide 30 translates this into NCCL terms: pull is Broadcast and push is Reduce.

Slide 46 names the problem: two synchronizations (parameters once, gradients once), with workers waiting for the server to send parameters and the server waiting for every worker's gradients.

NCCL's five building blocks

NCCL is NVIDIA's multi-GPU communication library. L14 slide 9 says it offers both collective and point-to-point communication over PCIe, NVLink, InfiniBand, and IP sockets, and ties every operation to a CUDA stream. Slides 10-15 introduce five collectives:

PrimitiveWhat it doesSlide
BroadcastCopy N elements from the root to all ranksp.11
ReduceReduce (sum, max, or min) across devices, write the result to one rankp.12
AllReduceReduce, and every rank gets the result (= Reduce + Broadcast)p.13
ReduceScatterReduce, then split the result into parts scattered across ranksp.14
AllGatherGather N values from each of k ranks into a k×N result sent to all ranksp.15

The bottom of slide 15 has the key identity: AllReduce = ReduceScatter + AllGather. Ring AllReduce is built exactly this way.

Why a ring: chop the data so transfers pipeline

Slide 18 says NCCL moves data and performs reductions around rings. Slides 19-28 derive why using Broadcast. Send N bytes at bandwidth B across K GPUs in a one-way ring:

  • Send it whole: each hop takes N/B, and there are K−1 hops, so total time is (K−1)·N/B (p.22). More GPUs means slower.
  • Split into S messages: each hop takes N/(S·B), and the second message can leave as soon as the first reaches the next GPU. Total time is (K−2+S)·N/(S·B), which is about N/B when S is large (p.28).

The second is nearly independent of GPU count. Chopping big messages so every link stays busy at once is the heart of the ring algorithm.

Ring AllReduce: ReduceScatter, then AllGather

Slides 34-43 draw ring AllReduce step by step with 4 workers. Each worker splits its gradient into 4 chunks (a0-a3 on worker 0, b0-b3 on worker 1, and so on):

  1. ReduceScatter phase (p.35-42): at each step, every worker sends one chunk to the next worker in the ring, which adds it to its own matching chunk. After K−1 steps, each worker holds one fully reduced chunk (for example, a1+b1+c1+d1).
  2. AllGather phase (p.43): another K−1 steps pass the finished chunks to everyone. At the end, every worker has the full summed gradient.

Slides 44-45 give MPI code for both phases. The slides don't derive the total communication volume of ring AllReduce; for that, CS336 Lecture 7 works through the math.

Back to the comparison on slide 46: AllReduce data parallelism needs no server, but every worker updates the parameters itself. The slide asks "redundant?" and answers that a local update is much faster than moving data between GPUs.

The NCCL call example on L14 (p.29)

The slide shows how one process managing several GPUs initializes communicators, launches an AllReduce, and waits for the streams:

NCCLCHECK(ncclGroupStart());
for (int i=0; i<nDev; i++) {
  CUDACHECK(cudaSetDevice(localRank*nDev + i));
  NCCLCHECK(ncclCommInitRank(comms+i, nRanks*nDev, id, myRank*nDev + i));
}
NCCLCHECK(ncclGroupEnd());

NCCLCHECK(ncclGroupStart());
for (int i=0; i<nDev; i++)
  NCCLCHECK(ncclAllReduce((const void*)sendbuff[i], (void*)recvbuff[i],
                          size, ncclFloat, ncclSum, comms[i], s[i]));
NCCLCHECK(ncclGroupEnd());

for (int i=0; i<nDev; i++)
  CUDACHECK(cudaStreamSynchronize(s[i]));

Mechanism 2: how PyTorch DDP hides the synchronization

L15 picks up from L14: the algorithm exists, so how does a real framework use it? The subject is PyTorch Distributed Data Parallel (DDP).

Two design goals

L15 slide 8 quotes the paper's two goals:

  • Non-intrusive: developers should be able to reuse their local training script with minimal changes
  • Interceptive: the API must let the implementation intercept signals and trigger the right algorithms promptly, exposing as many optimization opportunities as possible

The example on slide 12 shows the first goal in practice: wrap model as DDP(model, device_ids), and the loss, optimizer, backward(), and step() stay the same.

World size, global rank, local rank

Slide 9 defines three terms, using two machines with two GPUs each:

  • world size: total number of processes (4 here)
  • global rank: the process's global ID (0-3)
  • local rank: its ID within one machine (0 and 1 on each)

Slides 10-11 explain that the older torch/distributed/launch.py passes these through the environment variables MASTER_ADDR, MASTER_PORT, RANK, and WORLD_SIZE, and each process joins the process group with dist.init_process_group(backend="nccl"). The course's code example (the code walkthrough on L15 slide 23) launches with torchrun --nproc_per_node=2 main.py --ddp instead.

When to sync: gradient bucketing

The obvious approach is to wait for the whole backward pass, then AllReduce every gradient once. Slides 13-14 point out you can do better: gradients are computed layer by layer from the back, so the last layers' gradients are ready long before the first layers'. There's no reason to wait.

But syncing each parameter the moment it's ready syncs too often. DDP's compromise is buckets (slides 15-18):

  • Bucket size is set by the bucket_cap_mb argument to the DDP constructor
  • The mapping from parameters to buckets is fixed when DDP is constructed, based on the size limit and parameter sizes
  • Parameters go into buckets in roughly the reverse order of model.parameters(), because DDP expects gradients to become ready in about that order during backward
  • Once every gradient in a bucket is ready, the Reducer launches an asynchronous AllReduce on it while backward keeps computing earlier layers

The result is backward computation overlapping with AllReduce communication. Slide 20 shows the skeleton of the Reducer's C++ code: after each parameter's gradient accumulates, autograd_hook decrements its bucket's pending count; when it hits zero, mark_bucket_ready sends ready buckets to AllReduce in bucket order. The charts on slides 21-22 show DDP's scalability and the latency reduction from overlapping communication.

What the paper adds beyond the slides

The DDP paper describes the PyTorch v1.5 implementation. Its abstract lists three acceleration techniques: gradient bucketing, overlapping computation with communication, and skipping gradient synchronization. When configured appropriately, the paper reports near-linear scalability on 256 GPUs.

Details worth noticing when you read it:

  • The default bucket is 25 MB. The paper calls bucket size the key trade-off: bigger buckets amortize communication overhead better, but each bucket waits longer for its gradients. It recommends measuring for your own use case; in its ResNet50 experiments on 16 GPUs with the NCCL backend, the fastest range was 10 MB to 25 MB.
  • Why buckets must go out in order. Each process rebuilds the autograd graph dynamically on every forward pass, so gradient ready order can differ between processes. If each process sent buckets as soon as they were ready, AllReduce contents would mismatch, producing wrong results or crashes. So all processes must use the same bucket order, and none can launch bucket i+1 before bucket i. The paper admits reverse model.parameters() order is only an approximation.
  • Unused parameters can hang backward. If some parameters don't take part in an iteration, their bucket never fills and backward hangs. DDP's find_unused_parameters option walks the autograd graph at the end of forward to mark them.
  • Skipping sync: the no_sync() context manager lets you skip synchronization during the small steps of gradient accumulation and sync once on the last step.
  • The paper supports three communication backends, NCCL, Gloo, and MPI, and compares NCCL with Gloo.

What these lectures don't cover

  • What if the model doesn't fit on one GPU: the next-lecture readings on L15's last slide are GPipe and Megatron-LM, that is, pipeline and tensor parallelism.
  • Every GPU storing a full optimizer state is wasteful: that's the ZeRO post. L10's last slide also lists the PyTorch FSDP paper as a reading, but the Spring 2026 syllabus attaches only the DDP paper to these two sessions.
  • Implementation: HW5 has you write data and pipeline parallelism yourself, on at least 2 GPUs.

How to read these lectures

  1. In L14, start with the ring broadcast derivation on slides 19-28. Work out (K−2+S)·N/(S·B) yourself to feel why chopping helps.
  2. Then read ring AllReduce on slides 34-43 with pen and paper, following all 4 workers until you know who sends which chunk to whom at each step.
  3. In L15, read slides 13-20 alongside Section 3 (System Design) of the DDP paper, focusing on the reasons behind bucketing and ordering.
  4. Finally, run the course's ddp_example and compare per-epoch time on one GPU versus two.

One thing you can do tonight: in any PyTorch training script, set bucket_cap_mb to 1, 25, and 100, run a few steps each, and watch how step time changes.

Further reading

References