Every distributed training step is a race between arithmetic and data movement. The arithmetic is easy to estimate: about six floating-point operations per parameter per token. The data movement is harder, because a modern GPU cluster is not one network but a ladder of very different links, from on-package HBM at terabytes per second down to the network card at tens of gigabytes per second. A parallelism strategy that puts its heaviest traffic on the wrong rung can run several times slower than the same strategy placed correctly.
This article walks the ladder rung by rung with checked datasheet figures, shows exactly how many bytes data parallelism, fully sharded data parallelism and tensor parallelism move per step, turns that into a small budget model you can run, and ends with how to measure the real numbers on your own cluster. Bisection bandwidth and oversubscription across the whole fabric are covered separately in Bandwidth per GPU (Bisection).
The bandwidth ladder
Datasheets mix units and directions, so normalise everything to bytes per second, per GPU, per direction before comparing. NVIDIA lists the H100 SXM's NVLink as 900 GB/s and PCIe Gen5 as 128 GB/s; both are bidirectional totals, so a single transfer sees about half. A ConnectX-7 port at 400 Gb/s is 50 GB/s per direction. The DGX H100 pairs each of its eight GPUs with one single-port 400 Gb/s ConnectX-7 for the compute fabric, which gives 3.2 Tb/s per node, and keeps separate Ethernet ports for storage and management.
| Rung | H100 SXM (per GPU, per direction) | B200 (per GPU, per direction) | Who uses it |
|---|---|---|---|
| HBM | 3.35 TB/s | 8 TB/s (64 TB/s across 8 GPUs) | every kernel; bounds decode and optimizer steps |
| NVLink via NVSwitch | 450 GB/s (900 GB/s total) | 900 GB/s (14.4 TB/s aggregate across 8) | tensor parallel, intra-node collectives |
| PCIe Gen5 x16 | 64 GB/s (128 GB/s total) | check your platform datasheet | host-to-device copies, NIC traffic without a direct path |
| Compute NIC | 50 GB/s (400 Gb/s ConnectX-7) | 50 GB/s (DGX B200 also ships 8x 400 Gb/s) | data parallel, pipeline, FSDP across nodes |
Two things stand out. First, the gap between NVLink and the NIC is about 9x on H100 and 18x on a B200 node that keeps 400 Gb/s NICs: faster GPUs widen the gap unless the network is upgraded too, which is why newer platforms move to faster NICs (Modal, for example, quotes 6,400 Gb/s of RDMA networking for B300 clusters against 3,200 Gb/s for other GPU types). Second, the NIC is only fast if traffic reaches it without crossing the CPU. GPUDirect RDMA lets the NIC read GPU memory through the shared PCIe switch; if the NIC sits on the other CPU socket, or the driver path is missing, traffic bounces through host memory and the inter-node rung quietly gets slower. See NVLink for the scale-up side.
Units, directions and bus bandwidth
Collective benchmarks report two numbers and confusing them is the commonest mistake in bandwidth discussions. Algorithm bandwidth is simply the buffer size divided by the time: algbw = S / t. Bus bandwidth rescales it by how many bytes each rank actually has to move, so the number can be compared with link speed regardless of how many ranks took part. The factors used by NVIDIA's nccl-tests are:
| Collective | busbw = algbw × | Bytes each rank sends for an S-byte buffer |
|---|---|---|
| AllReduce | 2(n-1)/n | 2(n-1)/n · S (reduce-scatter then all-gather) |
| ReduceScatter | (n-1)/n | (n-1)/n · S |
| AllGather | (n-1)/n | (n-1)/n · S, where S is the full gathered output |
| AlltoAll | (n-1)/n | (n-1)/n · S |
| Broadcast, Reduce | 1 | S |
The practical consequence: the best possible all-reduce time for an S-byte buffer on links of bandwidth B is t = 2(n-1)/n · S / B. For large n that is nearly 2S/B, and it does not shrink as you add GPUs. Bandwidth, not GPU count, sets the floor. For small messages a second term dominates: per-hop latency of a few microseconds multiplied by the number of steps, which is why NCCL switches between ring and tree algorithms by message size and why frameworks bucket many small gradients into tens of megabytes before reducing them. NCCL all-reduce covers the algorithms in detail.
What each parallelism strategy sends
Each parallelism strategy has a characteristic traffic pattern, size and frequency. With P parameters, b bytes per element, n ranks in the group, and a micro-batch of B sequences of length s with hidden size h over L layers:
| Strategy | Collective | Bytes per rank per step | Where it should run |
|---|---|---|---|
| Data parallel (DDP) | AllReduce of gradients | 2(n-1)/n · P · b | any rung; overlaps with backward |
| FSDP / ZeRO-3 | AllGather params twice, ReduceScatter grads | 3(n-1)/n · P · b | NVLink or a fast fabric; use hybrid sharding across nodes |
| Tensor parallel (Megatron) | AllReduce of activations, 2 forward + 2 backward per layer | ≈ 2(n-1)/n · 4 L B s h b per micro-batch | inside one NVLink domain only |
| Pipeline parallel | point-to-point activations at stage boundaries | B s h b per boundary per micro-batch | across nodes is fine |
| Expert parallel (MoE) | AlltoAll of tokens, twice per MoE layer | depends on routing; bisection-bound | see bisection article |
Two features matter as much as the totals. Data-parallel traffic happens once per step and can overlap with the backward pass, so it tolerates the slow rung. Tensor-parallel traffic happens four times per layer per micro-batch and sits on the critical path between matrix multiplies, so it cannot hide. That is the whole argument for the usual layout: tensor parallel inside the node, pipeline and data parallel across nodes. Tensor parallelism and FSDP cover the mechanics.
A communication budget model
The model below turns those formulas into milliseconds. It uses ring lower bounds, a flat 80% link efficiency, and a hierarchical all-reduce of the kind NCCL builds on rail-optimised clusters: reduce-scatter over NVLink so each GPU owns one eighth of the buffer, all-reduce that shard across nodes over the GPU's own NIC, then all-gather back over NVLink. It ignores latency terms and congestion, so treat every result as a ceiling on what the hardware allows, not a prediction.
GB = 1e9
HW = { # per GPU, bytes/s per direction (datasheet peaks)
"h100": dict(nvlink=450 * GB, nic=50 * GB),
"b200": dict(nvlink=900 * GB, nic=50 * GB),
}
ALLREDUCE = lambda n: 2 * (n - 1) / n
GATHER_OR_SCATTER = lambda n: (n - 1) / n
def ring_time(nbytes, n, bw, factor):
return 0.0 if n == 1 else factor(n) * nbytes / bw
def hier_allreduce(nbytes, g, nodes, hw, eff=0.8):
"""Reduce-scatter inside the node, all-reduce each 1/g shard across
nodes on that GPU's own NIC (its rail), then all-gather inside the node."""
nvl, nic = hw["nvlink"] * eff, hw["nic"] * eff
intra = 2 * ring_time(nbytes, g, nvl, GATHER_OR_SCATTER)
inter = ring_time(nbytes / g, nodes, nic, ALLREDUCE)
return intra, inter
def fsdp_bytes(params, n, b=2):
return 3 * (n - 1) / n * params * b
def tp_bytes(batch, seq, hidden, layers, b=2):
return 4 * layers * batch * seq * hidden * b # payload, before the ring factor
def compute_s(params, tokens, gpus, peak=989e12, mfu=0.4):
return 6 * params * tokens / (gpus * peak * mfu)The 989 TFLOP/s peak is H100 SXM dense BF16 (NVIDIA's headline 1,979 figure assumes 2:4 sparsity), and 40% model FLOPs utilisation is a reasonable planning number for a well-tuned dense model. Swap in your own measured bandwidths once you have them; the structure is what matters.
Worked example: 8B on 128 H100s
Take an 8-billion-parameter dense model on 16 H100 nodes (128 GPUs), BF16 gradients, and a global batch of 4 million tokens per step. Running the model gives:
- Compute: 6 × 8e9 × 4e6 / (128 × 989e12 × 0.4) = 3.79 s per step.
- Gradient all-reduce (16 GB): intra-node phases 77.8 ms, inter-node phase 93.8 ms. Run back to back that is 171.5 ms; NCCL pipelines the phases, so the floor is close to the 93.8 ms inter-node part. Either way it is 2.5-4.5% of the step and overlaps with backward. Data parallel is not the problem here.
- The same all-reduce as a flat ring where every hop crosses a NIC: 793.8 ms. Hierarchy and rail placement are worth roughly 4.6x (phases back to back) to 8.5x (fully pipelined) on this one collective.
- On B200 with the same 400 Gb/s NICs: intra-node time halves to 38.9 ms, the inter-node 93.8 ms does not move at all, and compute gets roughly twice as fast. The network becomes a larger share of the step.
- Tensor parallel, TP = 8, one micro-batch of 8,192 tokens, h = 4,096, 32 layers: 8.59 GB of all-reduce payload. Over NVLink that is 41.8 ms against 124 ms of compute for the micro-batch: noticeable but tolerable. Over the NIC it would be 375.8 ms, three times the compute, which is why TP never leaves the node.
- FSDP of a 70B model sharded flat across 128 GPUs: 416.7 GB moved per GPU per step, 10.4 s at 40 GB/s. At 4M tokens the step computes for 33.2 s and prefetching can hide it; at 1M tokens compute drops to 8.3 s and the job becomes network-bound. Hybrid sharding (shard inside the node, replicate across nodes) moves most of those bytes onto NVLink.
The general lesson: compute time scales with tokens, communication for DP and FSDP scales with parameters. Shrinking the batch, or adding GPUs at a fixed batch, pushes you toward the network ceiling. Run the model before changing either.
Measuring the real numbers
Datasheet numbers are ceilings; the cluster you rent has cables, firmware, switch buffers and placement of its own. Measure each rung separately, then the collectives your job actually uses.
# GPU <-> GPU, NIC and CPU affinity matrix (NV#, PIX, PXB, NODE, SYS)
nvidia-smi topo -m
# Host <-> device and device <-> device copy bandwidth (github.com/NVIDIA/nvbandwidth);
# run with no arguments for the full test suite
./nvbandwidth
# Raw NIC throughput between two hosts (perftest); server first, then client
ib_write_bw -d mlx5_0 --report_gbits
ib_write_bw -d mlx5_0 --report_gbits <server-host>
# Collectives (github.com/NVIDIA/nccl-tests): one node, 8 GPUs, 8 B .. 8 GB doubling
./build/all_reduce_perf -b 8 -e 8G -f 2 -g 8
# 16 nodes, one process per GPU (build nccl-tests with MPI=1)
# (Open MPI: -x exports the variable to ranks on remote nodes)
mpirun -np 128 -N 8 -x NCCL_DEBUG=INFO ./build/all_reduce_perf -b 8 -e 8G -f 2 -g 1Read the busbw column at the largest sizes and compare it with the ceilings: a single node should approach the NVLink figure, and a multi-node all-reduce on eight 400 Gb/s rails is ultimately bounded by the per-node NIC total. With NCCL_DEBUG=INFO the log states which transport each channel chose. NET/IB with GPUDirect is what you want across nodes; NET/Socket means NCCL fell back to TCP and your inter-node rung has collapsed. Record the results per node pair and keep them: the same test after a firmware upgrade or cable swap is how regressions are caught.
Failure modes
- Bits versus bytes. 400 Gb/s is 50 GB/s. A capacity plan that mixes them is off by 8x.
- Bidirectional totals. 900 GB/s of NVLink is 450 GB/s each way; plans built on the total are 2x optimistic.
- Slowest link wins. A ring runs at the speed of its worst hop. One NIC negotiated at a lower speed, or one flapping cable, drags the whole job; per-node nccl-tests sweeps find it quickly.
- NIC on the wrong socket. If
nvidia-smi topo -mshows SYS between a GPU and its NIC, traffic crosses the CPU interconnect. Fix the mapping or the container's device assignment. - Silent TCP fallback. Missing RDMA libraries in the container make NCCL use sockets. The job still runs, slowly.
- Storage on the compute fabric. Checkpoint writes and dataset reads that share links with gradients cause periodic step-time spikes. Keep them on the front-end network or schedule them between steps. Rail-aligned topology explains the cross-rail case.
- Small-message storms. Unbucketed gradients or per-layer tiny collectives are latency-bound; they never reach link speed whatever the hardware.
Trade-offs and fixes
Bandwidth problems have four kinds of fix, in roughly increasing cost. Placement: keep tensor parallel inside NVLink, align data-parallel ranks with rails, and pack jobs onto whole nodes. Overlap: bucket gradients and start reducing during backward, prefetch FSDP parameters one layer ahead. Fewer bytes: hybrid sharding, BF16 or FP8 communication where the optimizer tolerates it, larger batches so compute grows relative to traffic. More bandwidth: more or faster NICs, which is a hardware purchase. Most teams that think they need the last option have not finished the first three.
What to do next
- Write down the four rungs for your nodes in bytes per second per direction, from datasheets.
- Run
nvidia-smi topo -mand confirm every GPU has a NIC on its own PCIe switch. - Run nccl-tests all_reduce_perf on one node and across two nodes; compare busbw with your ceilings.
- Plug your model size, batch and GPU count into the budget model and find which collective is closest to the step time.
- Check that tensor-parallel groups never span nodes and that NCCL logs show NET/IB, not NET/Socket.
- Re-run the sweep after every firmware, driver or cabling change and keep the results.