★ Pure arithmetic, like Section V.06. This is how you decide a parallelism strategy on paper before spending money on hardware.
1. The model#
Communication time for one collective:
T = α × hops + β × volume
α = per-hop latency NVLink ~2 µs, PCIe ~7 µs, IB ~3 µs
β = 1 / bandwidth NVLink ~1/450 GB/s, PCIe Gen5 ~1/64, IB ~1/50
hops = 2(N-1) for ring, 2·log₂(N) for tree
volume = 2 × (N-1)/N × message_size for ring AllReduceTwo regimes:
SMALL messages (< ~256 KB): latency-dominated. T ≈ α × hops
LARGE messages (> ~4 MB): bandwidth-dominated. T ≈ β × volumeLLM decode lives in the small-message regime. Prefill lives in the large-message regime. They need different reasoning.
2. The volumes, per strategy#
TENSOR PARALLELISM
per layer: 2 AllReduces (attention output, FFN output)
message_size = B × S × d × bytes
per token, per GPU:
volume = 2 × L × 2 × (N-1)/N × B × S_step × d × bytes
hops = 2 × L × 2(N-1) [ring] or 2 × L × 2log₂(N) [tree]
where S_step = 1 for decode, S_prompt for prefill
PIPELINE PARALLELISM
per stage boundary: one point-to-point send
volume = (P-1) × B × S_step × d × bytes
hops = P-1
EXPERT PARALLELISM
per MoE layer: 2 All-to-Alls
volume = 2 × L_moe × B × S_step × k × d × bytes
(All-to-All doesn't have the (N-1)/N reduction — each rank sends distinct
data to every other rank)
DATA PARALLELISM
volume = 03. Worked example 1 — TP decode#
Llama-3-70B, L=80, d=8192, BF16, TP=8, batch 32, decode (S_step=1)
message_size = 32 × 1 × 8192 × 2 = 512 KB
Per AllReduce, per GPU:
volume = 2 × (7/8) × 512 KB = 896 KB
ring hops = 14
Per token: 160 AllReduces (2 per layer × 80 layers)
total volume = 160 × 896 KB = 143 MB
total hops = 160 × 14 = 2,240
NVLink (α=2 µs, β=1/450 GB/s):
latency term: 2,240 × 2 µs = 4.48 ms ← DOMINATES
bandwidth term: 143 MB / 450 GB/s = 0.32 ms
T ≈ 4.8 ms
Hmm — that's as large as the compute time (5.3 ms)!
With TREE algorithm (NCCL uses this for small messages):
hops = 160 × 2 × log₂(8) = 160 × 6 = 960
latency term: 960 × 2 µs = 1.92 ms
T ≈ 2.2 ms → 42% overhead. Still significant.
With CUSTOM ONE-SHOT ALLREDUCE (vLLM/TRT-LLM):
effectively 1 hop per AllReduce (direct P2P writes)
latency: 160 × ~8 µs = 1.28 ms
T ≈ 1.6 ms → 30% overhead
MEASURED reality on H100 NVSwitch: TP=8 overhead is typically 15-25%,
because NCCL's LL protocol and the custom kernels do better than this
simple model, and some communication overlaps with compute.Two lessons from this example:
- The naive ring model badly overestimates because NCCL uses tree + LL for small messages.
- The latency term dominates decode. This is why TP scaling degrades with N and why custom all-reduce kernels matter.
4. Worked example 2 — TP prefill#
Same model, TP=8, prefill of a 2,048-token prompt, batch 1
message_size = 1 × 2048 × 8192 × 2 = 33.5 MB ← 65x larger than decode
Per AllReduce, per GPU:
volume = 2 × (7/8) × 33.5 MB = 58.7 MB
Per prefill: 160 AllReduces
total volume = 160 × 58.7 MB = 9.4 GB
NVLink:
bandwidth term: 9.4 GB / 450 GB/s = 20.9 ms ← now DOMINATES
latency term: 2,240 × 2 µs = 4.5 ms
T ≈ 25 ms
Compare to prefill compute:
FLOPs = 2 × 70e9 × 2048 = 287 TFLOP, /8 GPUs = 35.8 TFLOP each
at 600 TFLOP/s achieved: 60 ms
→ communication is 25/60 = 42% of compute. Significant but tolerable,
and much of it overlaps.Prefill is bandwidth-bound in communication; decode is latency-bound. Different regimes, different optimizations. Sequence-parallel norms and other techniques reduce prefill communication specifically.
5. Worked example 3 — PP vs TP across nodes#
Llama-3-405B, L=126, d=16384, FP8, 2 nodes × 8 H100
Question: TP=16 across nodes, or TP=8 × PP=2?
OPTION A: TP=16 (8 GPUs in each node, AllReduce crosses the node boundary)
message_size (decode, batch 32) = 32 × 1 × 16384 × 2 = 1 MB
per AllReduce per GPU: 2 × (15/16) × 1 MB = 1.875 MB
per token: 252 AllReduces × 1.875 MB = 472 MB
Half the ring hops cross InfiniBand:
IB portion: ~236 MB / 50 GB/s = 4.7 ms
plus IB latency: ~1,260 hops × 3 µs = 3.8 ms (half of 2×252×15)
T ≈ 8.5 ms of communication
Compute per GPU: 405e9 / 16 × 1 byte = 25.3 GB / 3350 GB/s = 7.6 ms
→ 112% overhead. UNUSABLE.
OPTION B: TP=8 within each node, PP=2 across
TP communication: stays on NVLink. As example 1, ~1.6-2 ms.
PP communication: 1 activation send per token
volume = 32 × 1 × 16384 × 2 = 1 MB
over IB: 1 MB / 50 GB/s = 0.02 ms + 3 µs latency
T ≈ 2 ms of communication
Compute per GPU: 405e9 / 8 × 1 byte = 50.6 GB / 3350 = 15.1 ms
→ 13% overhead. GOOD.
Plus pipeline bubbles: with 32 microbatches and P=2, (2-1)/(32+1) = 3%.
ANSWER: TP=8 × PP=2. Option A's cross-node AllReduce makes it unusable.This example is the justification for the canonical layout, derived rather than asserted.
6. Worked example 4 — when does TP stop helping?#
Model with weight bytes W, on GPUs with bandwidth BW, at TP=N:
compute time = W / (N × BW) ← shrinks as 1/N
comm time ≈ α × c × L × hops(N) ← grows as N (ring) or log N (tree)
total(N) = W/(N × BW) + α × c × L × 2log₂(N)
Llama-3-70B, W=141 GB, BW=3.35 TB/s, α=2 µs, L=80, c=2 (AllReduces per layer):
N=1: 42.1 ms + 0 = 42.1 ms
N=2: 21.0 ms + 0.64 ms = 21.6 ms (1.95x)
N=4: 10.5 ms + 1.28 ms = 11.8 ms (3.57x)
N=8: 5.3 ms + 1.92 ms = 7.2 ms (5.85x)
N=16: 2.6 ms + 2.56 ms = 5.2 ms (8.10x)
N=32: 1.3 ms + 3.20 ms = 4.5 ms (9.36x)
N=64: 0.7 ms + 3.84 ms = 4.5 ms (9.36x) ← no further gain
Efficiency (speedup/N):
N=2: 98% N=4: 89% N=8: 73% N=16: 51% N=32: 29% N=64: 15%The optimal N is where the marginal compute saving equals the marginal communication cost. Setting the derivative to zero:
d/dN [W/(N·BW)] = -W/(N²·BW)
d/dN [α·c·L·2·log₂N] = 2·α·c·L/(N·ln2)
Equal when: N = W·ln2 / (2·α·c·L·BW)
= 141e9 × 0.693 / (2 × 2e-6 × 2 × 80 × 3.35e12)
= 9.77e10 / 2.14e9 = 45.6So the time-optimal TP degree is ~46 — but efficiency there is 20%, meaning you’re using 46 GPUs to get 9x. Time-optimal is not cost-optimal.
Cost-optimal: the largest N where efficiency stays above your threshold.
Threshold 80% → N = 4-8
Threshold 70% → N = 8
Threshold 50% → N = 16Practical answer: TP ≤ 8, and use DP for throughput. Which is what everyone does — now you can derive why.
7. The estimation checklist#
For any proposed distributed configuration:
1. Compute per-GPU weight bytes and the resulting compute time.
2. Compute the communication volume and hop count per token.
3. Look up α and β for the actual interconnect (measure them — file 08).
4. T_comm = α × hops + β × volume.
5. Overhead = T_comm / T_compute.
6. If overhead > 30%, reconsider the strategy.
7. Check: does the TP degree divide n_kv_heads?
8. Check: are all TP-group GPUs in one NVLink domain?
9. Compute efficiency = speedup / N. If < 70%, use a smaller N + DP.
10. Compute the pipeline bubble if using PP: (P-1)/(M+P-1).8. Production implications#
- Do this arithmetic before provisioning. It takes fifteen minutes and can save six figures.
- Measure α and β on your actual hardware. Spec-sheet numbers are optimistic.
- Communication overlaps with compute to some degree in good implementations — expect your measured overhead to be 30-50% lower than this model predicts. Use the model to rank options, not to predict absolutes.
- Re-derive when the batch size changes. Decode communication scales with batch; at batch 1 you’re pure latency, at batch 512 you’re bandwidth-bound.
- The custom all-reduce matters most at low batch. Verify it’s enabled.
9. Common mistakes#
Using ring hop counts when NCCL uses trees. Overestimates small-message cost by 2-3x.
Ignoring the latency term for decode. It dominates.
Ignoring the bandwidth term for prefill. It dominates there.
Optimizing for time instead of efficiency. Time-optimal TP wastes GPUs.
Using spec-sheet bandwidth. Measure it.
Forgetting that communication partially overlaps compute.
Not re-deriving for a different batch size.
10. Hands-on exercise#
A. Build the calculator. Write a script that takes model config, GPU spec, interconnect parameters, parallelism strategy, and batch size, and outputs: compute time, communication time, overhead, efficiency, and a recommendation. This is the companion to your Section V.15 calculator.
B. Measure α and β. Use nccl-tests to fit α and β for AllReduce on your hardware (the
intercept and slope of time vs message size). Compare to the spec sheet.
C. Validate the model. For a model you can run at TP=2, 4, 8, predict the communication overhead with your calculator and compare to the measured difference in step time. What’s your model’s error? Calibrate it.
D. Derive the optimum. For three models and your hardware, compute the efficiency curve vs TP degree. Where does each drop below 70%? Does the answer match what your engine’s documentation recommends?
E. Multi-node decision. For a model too large for one node, compute the overhead for TP=16 vs TP=8×PP=2 using your measured IB numbers. Confirm example 3’s conclusion on your hardware.
11. Interview questions#
- Write the communication cost model. What are α and β?
- Why is decode latency-bound and prefill bandwidth-bound in communication?
- Compute the TP=8 communication overhead for a 70B model at batch 32 on NVLink.
- Why does TP scaling efficiency degrade, and what’s the functional form?
- Why is the time-optimal TP degree not the cost-optimal one?
- Derive why TP=16 across nodes is worse than TP=8×PP=2.
- How would you validate a communication cost model against measurements?
12. Further reading#
- [ESTABLISHED] Pope et al., “Efficiently Scaling Transformer Inference” (2022) — the analytical approach, done thoroughly
- [FUNDAMENTAL] Hockney’s α-β model for collective communication
- [REFERENCE] NCCL tuning documentation and
nccl-tests - Next: 10 — Multi-node inference