Systems · Chapter 6 of 7 · Why it’s built this way

Why the network looks this way

A big AI model is too large for one chip, so it’s split across many. How you split it decides how much the chips must talk to each other.

Data parallelism copies the model and splits the data. Tensor, pipeline and expert parallelism split the model itself. Each creates different traffic, from huge group exchanges to small hand-offs between neighbors.

Data, tensor, pipeline, sequence and expert parallelism, their communication volumes and collective patterns, overlapping communication with compute, and why tensor parallelism stays inside the scale-up domain while data parallelism crosses the scale-out network.

An AI model is a huge list of numbers. Training it means showing it examples, checking how wrong it was, and nudging the numbers to do a bit better. Then it does that again, hundreds of thousands of times.

The biggest models are trained on thousands of chips at once. Meta used about 16,000 chips for one of its models. No single chip could hold a model that big, let alone train it fast enough. So the work has to be split up.

Chips work the same way. Small groups of chips sit close together, joined by very fast links (see scale-up). The groups are joined by a slower network that spans the whole building (see scale-out). How you split the work decides which talk goes where.

Training a large language model means running a hundreds of thousands of times: push a batch of examples forward through the model, work out how wrong it was, push corrections backward, and update every weight. For a model with hundreds of billions of , one accelerator has neither the memory to hold all of that nor the speed to finish in a reasonable time. Llama 3’s largest model, for example, was trained on up to 16,000 GPUs with 80 GB of memory each.

There are four basic ways to divide the work, and real systems combine them:

  • copies the model and splits each batch of examples.
  • splits the math inside each layer.
  • splits the stack of layers into stages.
  • spreads the “experts” of a model across accelerators.

Each strategy makes the accelerators exchange different data, in different amounts, at different moments. Some exchanges are large and urgent; others are small or can wait. This chapter explains each strategy, measures its traffic, and shows why it explains the hardware: a small of very fast links (the scale-up chapter) and a large, slower (the scale-out chapter).

You know a training cluster is thousands of accelerators joined by a network. This chapter is about why that network is built in two very different tiers. The reason is the parallelization strategy. A training job is a grid of parallelism degrees, data (DP), tensor (TP), pipeline (PP), expert (EP) and sequence or context, whose product is the accelerator count. Each dimension has a characteristic collective, volume per step and latency tolerance:

  • TP issues several all-reduces per layer per microbatch, on the critical path.
  • EP issues all-to-alls per MoE layer.
  • PP sends point-to-point activations at stage boundaries.
  • DP all-reduces (or reduce-scatters) gradients once per step, overlapped with the backward pass.

The published recipe, from Megatron-LM onward, is to keep TP inside a server-sized scale-up domain and let PP and DP span the scale-out fabric. Llama 3 ordered its dimensions [TP, CP, PP, DP] from innermost to outermost for exactly this reason. We’ll derive the memory and traffic of each dimension with the standard open models (ZeRO, Megatron, GPipe, Korthikanti et al.), map them onto scale-up and scale-out links, and look at how real clusters and network designs follow from them. The collective libraries, compilers and kernels that implement all this belong to the Software guide; here they appear only as traffic.

The model1234layerreplicated ×4batch of examplessplit across chipsChip 1allChip 2allChip 3allChip 4allcommunicationonce per step
Strategy

Data parallelism: full replica on every accelerator, a quarter of the batch each. One gradient all-reduce per step.

The four basic strategies on four accelerators. Colors show which accelerator holds each piece of the model (4 layers × 4 slices) and which examples it processes.Share freely with credit: ‘Figure from chipfieldguide.com’

The simplest plan is to give every chip a full copy of the model. Each chip learns from different examples, all at the same time. This is called .

After each round, every chip has its own list of nudges for the model’s numbers. If each chip used only its own list, the copies would slowly drift apart. So the chips average their lists, and every copy makes the same change. This group swap is called an .

The swap happens only once per round. Chips can even start sending while they are still working. So this plan doesn’t mind a slower network.

The catch is memory. Every chip must hold the whole model, plus notes about each number. For a big model, that is far more than one chip can store. One fix is for the copies to split up the notes, so each chip keeps only its share.

In , dd accelerators each hold a full replica of the model. A batch of, say, 1,024 examples is split so each replica processes 1,024/d1{,}024/d of them. Every replica runs the forward and backward pass on its share and ends up with its own : one number per weight saying how to nudge it.

Before updating the weights, the replicas must agree, so they average their gradients with an , a operation in which every participant ends with the sum of everyone’s inputs. Then each applies the same update and the replicas stay identical.

Two properties make data parallelism easy on the network:

  • It talks once per step. If a replica processes its share in several pieces, it accumulates the gradients locally and synchronizes once at the end.
  • It can overlap. Gradients for the last layers are ready first. PyTorch’s DistributedDataParallel groups them into buckets and starts the all-reduce on each bucket while the backward pass is still computing earlier layers.

The problem is memory. Training with the popular Adam optimizer in 16-bit “mixed precision” stores, per weight, 2 bytes of weight, 2 bytes of gradient and 12 bytes of (a 32-bit master copy and two running averages): 16 bytes per parameter. A 1.5-billion-parameter model already needs at least 24 GB, against 3 GB for its 16-bit weights alone. Plain data parallelism stores all of it on every replica.

(Zero Redundancy Optimizer) removes that duplication in three stages: stage 1 splits the optimizer state across the dd replicas, stage 2 also splits the gradients, and stage 3 splits the weights too, gathering each layer’s weights from the others just before using it. PyTorch’s Fully Sharded Data Parallel (FSDP), a ZeRO-inspired design, also shards all three. In the ZeRO paper’s example, a 7.5-billion-parameter model needs 120 GB per device with plain data parallelism and 1.9 GB of model state with all three stages across 64 devices.

DP of degree dd replicates the model and splits the global batch BB. Gradients are synchronized once per step, after local accumulation over microbatches, with an all-reduce that a bandwidth-optimal ring performs by sending 2(d−1)/d2(d-1)/d of the gradient buffer per rank: a reduce-scatter followed by an all-gather. For Ψ\Psi parameters in bf16 that is about 4Ψ4\Psi bytes per rank per step, independent of dd for large dd, which is why DP scales to thousands of replicas. Narayanan et al. note the ring time scales as 1−1/d1 - 1/d.

Overlap is the second reason DP tolerates the scale-out network. DDP buckets gradients in reverse registration order and launches an all-reduce per bucket as soon as its last gradient lands, so most of the exchange hides behind the remaining backward computation. The exposed part is roughly the last bucket plus whatever exceeds the backward time.

Memory is the limit. Mixed-precision Adam holds 16Ψ16\Psi bytes of model state per replica. ZeRO partitions it over the dd ranks:

StageShardedModel-state bytes per rankDP traffic per step
Plain DPnothing16Ψ16\Psi2Ψ2\Psi (all-reduce)
1 (PosP_{\mathrm{os}})optimizer state4Ψ+12Ψ/d4\Psi + 12\Psi/d2Ψ2\Psi
2 (Pos+gP_{\mathrm{os+g}})+ gradients2Ψ+14Ψ/d2\Psi + 14\Psi/d2Ψ2\Psi (reduce-scatter + all-gather)
3 (Pos+g+pP_{\mathrm{os+g+p}}, FSDP)+ weights16Ψ/d16\Psi/d3Ψ3\Psi (1.5×)

Volumes here are in elements, following the paper; multiply by 2 for bf16 bytes. Stage 3’s extra Ψ\Psi is the weight all-gather needed before both the forward and the backward pass. FSDP issues those all-gathers per wrapped unit, prefetching the next unit’s gather during the current unit’s compute, and discards the gathered weights afterwards. With gradient accumulation it can either reduce every microbatch (sharded gradients, more traffic) or keep unsharded gradients locally (more memory). ZeRO stages 1 and 2 are therefore nearly free on the wire, while stage 3 pays in latency-sensitive gathers that a fast network or good prefetching must hide.

Replica 1weight w10+3own updateReplica 2weight w10−1own updateReplica 3weight w10+2own updateReplica 4weight w10+4own updateall-reduce: mean of the updates = +2replicas identical
1 / 4

Four replicas hold identical weights. Follow one weight, w = 10.

One data-parallel step on four replicas. The all-reduce is what keeps the copies identical. Values are illustrative.Share freely with credit: ‘Figure from chipfieldguide.com’
Rep. 1Rep. 2Rep. 3Rep. 4weights (2 B)gradients (2 B)optimizer state (12 B)7.5B params, 64 devicesper device120 GBtraffic per step1×vs plain data parallel
ZeRO stage

Plain data parallelism: every replica stores all 16 bytes per parameter, 120 GB for 7.5B parameters.

ZeRO’s three stages, drawn for 4 replicas: solid is what each replica stores, dashed is held by the others. Numbers are the ZeRO paper’s example, 7.5B parameters on 64 devices, model state only.Share freely with credit: ‘Figure from chipfieldguide.com’

When a model still won’t fit, you split the model itself. There are three main ways.

Cut each layer into pieces

A model is built in layers, like the floors of a building. Each example passes through every layer in order. With , several chips share each layer and work on it at the same moment.

But they must add up their partial answers before the next layer can start, over and over. It’s like three friends adding a long column of numbers. Each takes a third, and they call out their totals before anyone can go on.

Line the layers up

With , chip 1 holds the first few layers, chip 2 the next few, and so on. It’s an assembly line. Each chip just hands its result to the next one, which is very little talking.

The downside is waiting. At the start, the later chips have nothing to do yet. At the end, the early ones sit idle. This idle time is called the . Sending lots of small pieces of work down the line keeps it busy.

Spread the specialists

Some models are made of many specialist parts, and each word visits only one or two of them. If the specialists live on different chips, each word must travel to its specialists and back. So every chip sends something to every other chip.

Data parallelism, even sharded, keeps every layer’s computation on one accelerator and stores all the intermediate results, the , of its examples. When those are too big, or there are more accelerators than the batch can feed, the model itself has to be split.

Tensor parallelism: split inside each layer

A transformer layer is dominated by matrix multiplications. Megatron-LM splits the first weight matrix of each block by columns and the second by rows across tt accelerators. Each accelerator computes a partial result, and the partial results are summed with an all-reduce. That takes two all-reduces in the forward pass and two in the backward pass of every layer. Weights, gradients and optimizer state are all divided by tt, and so is most of the computation.

The price is that each all-reduce sits between two steps that depend on it: the next operation cannot start until the sum arrives. That happens for every layer and every microbatch, so tensor parallelism needs the fastest, lowest-latency links available. A companion technique, Megatron , also splits the parts of the layer that tensor parallelism leaves whole (normalization and dropout) along the sequence of tokens. It cuts their activation memory by tt at no extra communication.

Pipeline parallelism: split the stack of layers

Pipeline parallelism gives each of pp accelerators a consecutive group of layers, a stage. To keep the stages busy, the batch is cut into mm that flow through the stages one after another. Between stages, only the activations at the boundary travel forward and their gradients travel back: point-to-point sends between two neighbors, not group operations.

The cost is the . Stage pp cannot start until the first microbatch has crossed stages 1 to p−1p-1, and stage 1 sits idle at the end while the last backward passes drain. The idle time equals p−1p-1 microbatch-times, compared with mm microbatch-times of real work. GPipe found the bubble negligible once there are at least four times as many microbatches as stages.

Expert parallelism: spread the experts

In a (MoE) model, each feed-forward layer is replaced by many experts, and a router sends each token to only a few. Parameters grow with the number of experts while the work per token stays almost constant. places the experts on different accelerators. Before each expert layer, every accelerator sends each token to the accelerators holding its chosen experts, and the results come back afterwards. Both moves are exchanges. Expert parallelism is usually laid over the data-parallel accelerators, which already process different tokens.

Context parallelism: split very long sequences

For very long inputs, splits the sequence itself across accelerators. Attention still needs every token to look at earlier ones, so the accelerators exchange keys and values, either passing them around a ring while computing or gathering them all at once. Llama 3 used 16-way context parallelism to train on 131,072-token sequences.

TP: column-then-row splits

For the MLP Y=GeLU(XA)Y = \mathrm{GeLU}(XA), Z=YBZ = YB, Megatron splits AA by columns and BB by rows, so each of tt ranks computes GeLU(XAi) Bi\mathrm{GeLU}(XA_i)\,B_i with no communication in between. The partial ZiZ_i are summed by one all-reduce (operator gg) in the forward pass, and the conjugate operator ff all-reduces the input gradient in the backward pass. Attention splits QQ, KK, VV by heads the same way. That makes two all-reduces forward and two backward per layer. Each moves a b×s×hb \times s \times h tensor, so per layer per microbatch a rank sends 8bsh (t−1)/t8bsh\,(t-1)/t elements, and the total grows with layers per stage × microbatches.

Sequence parallelism replaces each all-reduce with a reduce-scatter and an all-gather along ss around the LayerNorm and dropout regions. Because a ring all-reduce is a reduce-scatter plus an all-gather, bandwidth is unchanged, and activation memory per layer drops from sbh (10+24/t+5as/(ht))sbh\,(10 + 24/t + 5as/(ht)) to sbh (34/t+5as/(ht))sbh\,(34/t + 5as/(ht)). TP also shrinks each GEMM by tt, which lowers arithmetic intensity, and Narayanan et al. flag both this and the inter-server all-reduce as the reasons TP degrades beyond one server.

PP: schedules and the bubble

With pp stages and mm microbatches per pipeline per step, GPipe runs all forwards, then all backwards. The bubble is (p−1)(tf+tb)(p-1)(t_{\mathrm{f}} + t_{\mathrm{b}}) against an ideal m(tf+tb)m(t_{\mathrm{f}} + t_{\mathrm{b}}), a fraction (p−1)/m(p-1)/m. GPipe must stash activations (or stage inputs, with recomputation) for all mm microbatches. The 1F1B schedule (PipeDream-Flush) has the same bubble but at most pp microbatches in flight. Interleaving vv model chunks per device cuts the bubble to (1/v)(p−1)/m(1/v)(p-1)/m at vv times the point-to-point traffic. Per microbatch, each stage boundary carries bshbsh elements in each direction. With TP, the scatter/gather optimization sends 1/t1/t of that tensor from each TP rank over its own NIC and re-gathers it over the scale-up fabric.

EP: dispatch and combine

With EE experts and top-kk routing, an MoE layer’s FLOPs track kk while its parameters track EE. Experts are sharded over an EP group usually carved from the DP dimension; Switch Transformer maps one expert per data-parallel core. Every MoE layer does a dispatch all-to-all and a combine all-to-all forward and the mirrored pair backward. Each sends about k bsh (EP−1)/EPk\,bsh\,(\mathrm{EP}-1)/\mathrm{EP} elements per rank, before capacity padding and imbalance. GShard measured all-to-all cost growing roughly as D\sqrt{D} for DD devices on its 2D TPU cluster; going from 128 to 2,048 experts raised its share of MoE-plus-Transformer time from 16% to 36%. Unlike ring collectives, all-to-all traffic crosses the bisection, so it is the dimension most sensitive to how the scale-out network is built.

Sequence and context parallelism

Megatron’s “sequence parallelism” (above) lives inside the TP group. Context parallelism (CP) instead partitions the whole sequence across ranks for every layer, including attention, and exchanges KK and VV. Ring Attention passes K/V blocks around a ring and hides each transfer behind blockwise compute while compute takes longer. Llama 3 instead all-gathers KK and VV, splitting the sequence into 2×CP2 \times \mathrm{CP} chunks for causal load balance. It reports the all-gather overhead as negligible, because attention compute grows as s2s^2 while the K/V exchange grows as ss (and grouped-query attention keeps KK and VV small).

×→×=XA, columnsY = GeLU(XA)B, rowsZ1 chip
1 / 5
TP degree t

One MLP block on a single accelerator: Z = GeLU(X·A)·B.

Megatron-style tensor parallelism of a transformer MLP block: column split, then row split, then one all-reduce. Shapes are schematic.Share freely with credit: ‘Figure from chipfieldguide.com’
time →bubble: 43%stage 1stage 2stage 3stage 411112222333344441121231234234344forwardbackwardbubble

Bubble: 3 of 7 slots on every stage, 43% of the step. GPipe found it negligible from m ≥ 4p.

A pipeline with flushes: forward passes (light) flow down the stages, backward passes (dark, about twice as long) flow back up. Gray is the bubble.Share freely with credit: ‘Figure from chipfieldguide.com’
Chip 1thecatsatonChip 2abigredhatChip 3weranupitChip 4mydogateallExpert 1Expert 2Expert 3Expert 4┊ = capacity
1 / 4
Router

Each accelerator holds four tokens and one expert. The router picks an expert for each token.

Expert parallelism with top-1 routing on four accelerators. Routing and the capacity of 5 tokens per expert are illustrative.Share freely with credit: ‘Figure from chipfieldguide.com’

Put the four plans side by side, and their talking looks very different. Splitting each layer is the chattiest by far, and the chips must wait for every swap. Specialists are busy too. Copies make one big swap per round, while they keep working. The assembly line is the quietest: each chip just passes work to its neighbor.

Talking while you work costs nothing. Talking while everyone waits is lost time. A big training run costs millions of dollars, so lost time is lost money.

Here is the traffic of each strategy for one accelerator during one training step, in the standard published models.

StrategyOperationHow oftenSizeCan it overlap?
Tensor (TP)all-reduce4 per layer, every microbatchone layer’s activations eachHardly: the next step needs the result
Expert (EP)all-to-all4 per MoE layer, every microbatchthe routed tokensOnly with careful scheduling
Pipeline (PP)send / receive2 per microbatchone layer’s activationsYes, mostly
Data (DP)all-reduce (or reduce-scatter + all-gather)once per stepall the gradients this accelerator holdsYes, behind the backward pass

A worked example, using the simulator’s defaults: a 175-billion-parameter model (96 layers) on 1,024 accelerators with 8-way tensor, 8-way pipeline and 16-way data parallelism, 2,048-token sequences and 1,024 sequences per step. Each accelerator holds 12 layers and processes 64 microbatches per step. Per step, it sends roughly:

  • TP: about 350 MB per layer per microbatch, or about 270 GB per step.
  • DP: its 2.7 billion gradients (5.4 GB in 16-bit), about 10 GB per step through the ring.
  • PP: under 1 GB.

Tensor parallelism outweighs everything else by more than an order of magnitude, and it is the one that cannot hide. That single fact explains most of the hardware.

Hiding communication behind computation

An accelerator can compute and communicate at the same time, so traffic only costs time if the computation has to wait for it.

  • Data parallelism overlaps naturally: gradient exchange for late layers runs while early layers are still in their backward pass.
  • FSDP fetches the next layer’s weights while the current layer computes.
  • Ring Attention passes chunks of the sequence around while computing on the previous chunk.
  • DeepSeek-V3, whose expert traffic took about as long as its computation, built a pipeline schedule that interleaves one microbatch’s communication with another’s computation.

Tensor parallelism is the hard case, because every exchange sits between two operations that depend on it.

Per rank per step, with b=1b = 1, mm microbatches per pipeline, L/pL/p layers per stage and bf16, the standard volumes are as follows.

DimCollectiveElements sent per rank per stepExposure
TP (+SP)4 all-reduce, or 4 RS + 4 AG, per layer8sh (t−1)/t⋅(L/p)⋅m8sh\,(t-1)/t \cdot (L/p) \cdot mcritical path; ×1.5 with full recompute (extra forward)
EP2 all-to-all fwd + 2 bwd per MoE layer4ksh (e−1)/e⋅(L/p)⋅m/t4ksh\,(e-1)/e \cdot (L/p) \cdot m/tcritical path unless overlapped across microbatches
PPsend/recv2sh⋅m/t2sh \cdot m/t (scatter/gather)mostly hidden; latency adds to the bubble
DPall-reduce or RS + AG2⋅Ψ/(tp)⋅(d−1)/d2 \cdot \Psi/(tp) \cdot (d-1)/d, ×1.5 for ZeRO-3overlaps the backward pass

Three observations follow.

  • TP and EP scale with tokens; DP scales with parameters. TP traffic per step is proportional to (tokens per replica)×h×(layers per stage)(\text{tokens per replica}) \times h \times (\text{layers per stage}). DP traffic is proportional to the parameter shard and independent of batch size. Growing the global batch therefore grows TP and PP traffic, while DP traffic per token falls.
  • Ring collectives need neighbor bandwidth; all-to-all needs bisection. A ring all-reduce of nn ranks finishes in about (S/B)⋅2(n−1)/n(S/B) \cdot 2(n-1)/n, where BB is the slowest link on the ring. An all-to-all pushes (n−1)/n(n-1)/n of every buffer across the cut between any two halves of the group.
  • Overlap differs by dimension. DP all-reduces hide behind the backward pass with bucketing. FSDP prefetches its gathers. PP sends and receives in both directions in parallel with compute. TP all-reduces separate dependent GEMMs and are largely exposed; Korthikanti et al. do overlap the extra all-gather that sequence parallelism adds in the backward pass. DeepSeek-V3 measured a computation-to-communication ratio of about 1:1 for its cross-node EP and designed DualPipe to overlap a forward chunk’s all-to-all with a backward chunk’s computation.

Narayanan et al. measured the result at scale. Training their 1-trillion-parameter model on 3,072 GPUs, the effective bisection bandwidth used was 892 GB/s for pipeline point-to-point traffic and 12.9 TB/s for data-parallel all-reduces.

sent per accelerator per stepTensor (TP)≈ 270 GBData (DP)≈ 10 GBPipeline (PP)< 1 GBExpert (EP)MoE models only (not in this example)computecommsone step (schematic), time →compute idle: 33%idle, waiting on comms
Strategy

TP: each all-reduce sits between two dependent matrix multiplies, so the compute waits. Overlap barely helps. Compute idle 33% of the time.

Top: data sent per accelerator per step in the chapter’s worked example (175B, TP 8 × PP 8 × DP 16). Bottom: schematic timelines; time units are illustrative.Share freely with credit: ‘Figure from chipfieldguide.com’

Now the network design makes sense. Chips come in boxes of about eight, joined by very fast, short links. The boxes are joined by a network that spans the whole building. That network reaches much farther, but it carries a few times less data per chip.

So almost every big training job follows the same recipe. The chattiest work, splitting each layer, stays inside one box. The assembly line and the copies run between boxes, because they talk less often or can talk in the background.

When Meta trained its Llama 3 model, each layer was shared by eight chips in one box. The copies of the model were spread across the building.

It works the other way too. The traffic is so predictable that designers can leave out network parts nobody uses. One study found this could cut the network’s cost by a third to three quarters, with no loss of speed.

A is a small group of accelerators, typically the eight in one server, joined by a dedicated high-bandwidth fabric. The of InfiniBand or Ethernet switches joins thousands of such domains. The gap per accelerator is large. In Megatron-LM’s original cluster, each server had 300 GB/s between its GPUs but 100 GB/s for the whole server to the network, and in DeepSeek-V3’s cluster NVLink offered 160 GB/s against 50 GB/s for InfiniBand.

Frameworks number the accelerators so that the most demanding parallelism uses neighbors. Tensor parallelism is “innermost”: ranks 0–7 form one tensor-parallel group, which lands inside one server. Data parallelism is outermost, so its groups stretch across servers.

  • Megatron’s authors summarized this as a rule: use tensor parallelism up to the number of GPUs in a server, then pipeline parallelism across servers, then data parallelism to scale further.
  • When the ZeRO team ran Megatron’s tensor parallelism for a 40-billion-parameter model across two servers, it managed about 5 TFLOPS per GPU, under 5% of peak.
  • Llama 3 chose its order, [TP, CP, PP, DP], so that the innermost dimension, which needs the most bandwidth and the lowest latency, stays inside a server. Data parallelism, with FSDP, is outermost because it can prefetch weights and reduce gradients in the background.

The influence also runs from the traffic back to the network design:

  • Rails. Once tensor parallelism stays inside a server, the traffic that leaves a server is mostly between accelerators with the same position in different servers (GPU 3 of server A to GPU 3 of server B). “Rail-optimized” networks connect each position to its own set of switches. A 2023 study found GPT-style training traffic of about 300 GB between pairs inside a scale-up domain against about 6 GB between pairs across the network. It proposed dropping the top switch layer altogether, cutting network cost by 38–77%.
  • Oversubscription. Meta’s 24,000-GPU cluster gives full bandwidth inside pods of 3,072 GPUs but only one-seventh of that between pods. Its parallelism layout and job scheduler try to keep traffic inside a pod.
  • Reshaping the topology. Google’s TPU v4 joins 4,096 chips with optical switches that can rewire the 3D torus for each job: a long thin shape such as 4×4×32 for pipelines, with data parallelism mapped along one dimension.

Rank layout is the bridge between the parallelism grid and the topology. Llama 3 numbers ranks in the order [TP, CP, PP, DP], innermost first. The simulator uses that order and slots an EP group, when present, just outside TP, so a rank’s index is r=ti+t (ei+e (pi+p di′))r = t_i + t\,(e_i + e\,(p_i + p\,d'_i)), with d′=d/ed' = d/e. A dimension with stride SS and size GG lies wholly inside a domain of UU ranks when SG≤USG \le U, and every hop crosses domains when S≥US \ge U. So:

  • TP stays inside when t≤Ut \le U.
  • EP stays inside when te≤Ute \le U.
  • PP neighbors cross domains once te≥Ute \ge U, which is almost always.
  • DP, outermost, spans the whole job and crosses whenever the job is larger than one domain.

The simulator applies exactly this rule and times a group that spans domains at the scale-out rate. A ring is only as fast as its slowest hop, unless a hierarchical algorithm is used.

Hierarchical collectives change the picture for DP. The Rail-only paper analyzes a two-phase all-gather over xx GPUs per domain and yy domains. Across the network it moves D(y−1)/xD(y-1)/x in total, all of it between same-rank GPUs (the same rail); inside each domain it moves D(x−1)D(x-1). For GPT-1T on 4,096 GPUs in 256-GPU domains, a hierarchical all-reduce put 98% of the DP traffic on the scale-up fabric. Combined with TP kept inside domains, the scale-out traffic of a dense model is sparse and rail-local. That is why the authors argue spine switches are unnecessary, at 38–77% lower network cost and only 8.2–11.2% all-to-all overhead for MoE, by forwarding through the scale-up domain.

MoE is where the scale-out network is stressed. DeepSeek-V3 used no TP at all, 16-way PP, 64-way EP spanning 8 nodes and ZeRO-1 DP on 2,048 H800 GPUs. To respect the 3.2× NVLink-to-InfiniBand gap (160 vs 50 GB/s), it limited each token to experts on at most 4 nodes. Each token crosses InfiniBand once per destination node, to the GPU with the same in-node index, and NVLink completes the fan-out, so a token can select an average of 3.2 experts per node without extra NVLink overhead. This is the same rail principle, applied to all-to-all by the model’s router rather than by the network.

Torus machines map dimensions onto physical axes instead of rank strides. TPU v4 places DP along one torus dimension and the two model-parallel factors on the others, and reports 1.2–2.3× gains from choosing topology and hyperparameters together.

Server 1scale-upServer 2scale-upServer 3scale-upServer 4scale-up0123456789101112131415scale-out network: one rail per GPU positionTP: in-server
Rank order
Show traffic

TP innermost: ranks 0–3 form a TP group inside server 1, and so on. All TP all-reduces stay on the fast links.

Rank layout on 4 servers × 4 GPUs (8 is typical), TP 4 × PP 2 × DP 2. Each GPU position has its own network rail. Colors mark the groups of the selected strategy.Share freely with credit: ‘Figure from chipfieldguide.com’

Pick a model size and a number of chips. Then choose how many chips share each layer, and how many assembly-line stages to use. The rest of the chips become copies of the model.

Watch the memory bar: if it passes the line, the model doesn’t fit. Then watch the talk bars. Green talk stays inside a fast box; red talk has to cross the slower building network. Try splitting each layer across more chips than one box holds, and see what turns red.

The planner computes memory per accelerator and traffic per step for a GPT-style model. The TP and PP sliders set those degrees, and DP is whatever is left: N÷(TP×PP)N \div (\mathrm{TP} \times \mathrm{PP}). Things to try:

  • Press “DP only” with the default 175B model. How far over memory is it? Now turn on ZeRO stages one by one.
  • Raise TP above the scale-up domain size. Which bar turns red, and how does the step time change?
  • Raise PP to 32 or 64 and lower the batch. Watch the bubble in the status bar.
  • Switch to 64 experts and raise EP. Where does the new all-to-all traffic go?

The Expert view adds activation recomputation, per-strategy collective types, time on the wire and an exposed-communication estimate. The model uses 12h212h^2 parameters per layer, b=1b = 1, bf16 with mixed-precision Adam (16 bytes per parameter before sharding), and Korthikanti activations with sequence parallelism. The first pipeline stage holds min⁡(p,m)×⌈L/p⌉\min(p, m) \times \lceil L/p \rceil layers of activations. Traffic follows the table above, and the bubble is (p−1)/m(p-1)/m. Step time is compute×(1+bubble)\text{compute} \times (1 + \text{bubble}) plus exposed communication; DP hides behind up to two-thirds of compute. Bandwidths (450 and 50 GB/s) and 400 TFLOP/s of sustained math are illustrative. Try:

  • Set the 406B model, 16,384 accelerators, TP 8, PP 16, ZeRO 3, a batch of 2,048 and an 8,192-token sequence, close to Llama 3’s Table 4. Compare selective and full recompute.
  • Set the domain to 4 with TP 8 and read the exposed communication.
  • With 64 experts, compare EP 8 (inside a domain of 8 when TP = 1) with EP 64.

Not modeled: interleaved schedules, embedding and logit layers, temporary buffers and fragmentation, MoE capacity padding and imbalance, hierarchical collectives, and latency. Treat the numbers as order-of-magnitude.

Loading simulation…
Chips used to train Llama 3
up to 16,000
Chips sharing each layer in that run
8, in one box
Bytes of memory per model number, while training
16
How much more data the in-box links carry (one cluster)
3.2×

What these numbers mean:

  • 16,000 chips worked on one model at the same time. When one breaks, the whole job may have to stop. Over 54 days, Meta’s run stopped unexpectedly 419 times, nearly eight times a day.
  • 8 chips in one box shared every layer, because that’s where the fastest links are.
  • 16 bytes per number is the true memory cost of training. (A byte holds about one letter of text.) It covers the number, its nudge and the learning notes. A model with 175 billion numbers needs about 2.8 trillion bytes, far more than any one chip holds.
  • 3.2× is how much more data the in-box links could carry each second than the network links, in one recent cluster. So engineers keep the chattiest talk inside the box.
Llama 3 405B, 16K H100: TP × PP × DP
8 × 16 × 128
Megatron 1T model, 3,072 A100: % of peak
52%
Training memory, mixed-precision Adam
16 B/param
DeepSeek-V3 NVLink vs InfiniBand
160 vs 50 GB/s

Reading the table:

  • Every dense run uses tensor parallelism of 8, small enough to fit inside one server. TP never left the scale-up domain.
  • Pipeline depth grew with model size (64 stages for the 1-trillion-parameter model), and data parallelism absorbed the remaining GPUs.
  • For long contexts, Llama 3 traded data parallelism for context parallelism: same GPUs, same TP and PP, but DP fell from 128 to 8 to make room for CP 16. MFU dropped a few points.
  • DeepSeek-V3, a mixture-of-experts model, dropped tensor parallelism and spent its communication budget on expert all-to-alls instead.

Points worth noting:

  • Utilization holds up at scale. Narayanan et al. saw superlinear weak scaling from 1B to 1T parameters, 44% to 52% of peak, because larger models have larger GEMMs without proportionally more communication. Llama 3 lost two points of MFU going from 8K to 16K GPUs, which Meta attributes to the batch per DP group halving.
  • PTD-P beat ZeRO-3 alone. For 175B and 530B models, Megatron’s combination outran ZeRO-3 without tensor parallelism by 70% when GPUs were doubled at constant batch, “due to less cross-node communication.”
  • Bandwidth ratios set the rules. Megatron-LM’s 2019 cluster had 300 GB/s between GPUs inside a server but 100 GB/s for the whole server to the network, shared by all its GPUs. DeepSeek-V3’s cluster had a 3.2× NVLink-to-InfiniBand gap. A smaller gap makes cross-node EP viable, while a larger one pins more dimensions inside the domain.
  • Failures scale with NN. Meta logged 466 job interruptions in 54 days of 16K-GPU training, 78% of the unexpected ones hardware-related, and still reached over 90% effective training time.

Every way of splitting the work trades one problem for another.

  • Copies of the model are simple and patient. But each copy needs room for the whole model, unless the copies share out their notes.
  • Splitting each layer saves memory, but it needs the fastest links. So it only works inside a small group.
  • An assembly line talks little. But stations wait at the start and end of each round, and the first station must remember a lot of half-done work.
  • Specialists give a much bigger model for the same work per word. But every word must travel. And if too many words want the same specialist, some chips are swamped while others sit idle.

Size has a hidden cost too. When thousands of chips work in lockstep, one failure can stop them all. The bigger the job, the more often something breaks.

StrategyGainsCostsTypical limit
Data (plain)Simple; traffic once per step and hiddenFull model state on every replicaBatch size: each replica needs examples
Data + ZeRO/FSDPModel state divided by ddStage 3: 1.5× traffic, latency-sensitive gathersNetwork latency for the gathers
TensorDivides weights, compute and (with SP) activations4 exposed all-reduces per layer; smaller, less efficient matrix mathThe scale-up domain size
PipelineDivides weights; cheap point-to-point trafficBubble; first stage holds activations for pp microbatchesMicrobatches per step (m≫pm \gg p)
ExpertMany more parameters per FLOPAll-to-all traffic; load imbalanceBisection bandwidth

Some ways a configuration goes wrong:

  • Tensor parallelism spilling out of the server. All-reduces then run at network speed and the accelerators idle; the ZeRO paper measured under 5% of peak in this case.
  • Too few microbatches. As data parallelism grows at a fixed global batch, each pipeline gets fewer microbatches and the bubble grows. Llama 3’s utilization fell from 43% to 41% when it went from 8K to 16K GPUs, which Meta attributes to the smaller batch per data-parallel group.
  • Activations, not weights, run out. Pipelines divide the weights but not the first stage’s activations, so long sequences can overflow memory even when the weights are tiny. The usual fixes are recomputation or sequence parallelism.
  • Expert hot spots. If the router favors a few experts, their accelerators become the bottleneck. MoE systems add load-balancing terms or capacity limits.
  • TP degree vs GEMM efficiency. Every doubling of tt halves each rank’s GEMM width and doubles the relative cost of its all-reduces. Beyond t=gt = g (GPUs per server) the all-reduce also moves to slower links. Narayanan et al. measured up to 2× lower throughput for poor TP/PP splits even with fast inter-server links.
  • PP bubble vs memory vs batch. The bubble (p−1)/m(p-1)/m falls with mm, but m=B/(db)m = B/(db), so at a fixed global batch (set by optimization, not hardware) more DP means fewer microbatches. Interleaving trades bubble for v×v\times point-to-point traffic. Llama 3 had to modify its schedule because existing implementations required the per-GPU batch to be divisible by the number of stages.
  • Microbatch size. Larger bb raises kernel arithmetic intensity but lowers mm, which enlarges the bubble. The optimum is model- and configuration-dependent; Narayanan et al. found up to 15% throughput at stake.
  • Recomputation. Full recomputation costs one extra forward pass, up to 33% throughput at small batch. Yet it enabled batch sizes that gave up to 2× higher throughput through a smaller bubble. Selective recomputation with sequence parallelism cut activation memory 5× while removing over 90% of the recompute overhead.
  • ZeRO-3 vs model parallelism. FSDP avoids model-code changes but its gathers grow with re-gathering per microbatch and are latency-bound at small shard sizes. Llama 3 keeps it outermost and relies on prefetch; Megatron found PTD-P 70% faster than ZeRO-3 alone at 175B and 530B.
  • EP vs network topology. All-to-all is the one pattern that wants bisection bandwidth. Rail-only networks absorb it by forwarding through the scale-up domain at 8.2–11.2% overhead, DeepSeek-V3 by node-limited routing. The model’s routing and the network’s tiers end up co-designed.
memory savedtrafficidle timeData (plain)Data + ZeRO-3/FSDPTensorPipelineExpert

Tap or hover a row for its gains, costs and limit.

The strategies’ trade-offs at a glance. Dot ratings (0–3) are a qualitative summary of the table above, not measurements.Share freely with credit: ‘Figure from chipfieldguide.com’
□ = 256 chipsone week (rows = days); ✕ = interruptionday 1day 2day 3day 4day 5day 6day 7≈ 7.8 interruptions/day

16,384 accelerators: ≈ 7.8 unexpected interruptions per day, one every 3.1 hours, if the per-chip rate stays at Llama 3’s.

Interruptions grow with the number of accelerators. Scaled from Llama 3’s 419 unexpected interruptions in 54 days on 16,384 GPUs, assuming a constant per-chip rate. Marks are evenly spaced averages.Share freely with credit: ‘Figure from chipfieldguide.com’

This part goes deeper, into the math, models and algorithms behind the chapter. It’s written for the Expert level.

The planner’s model, assembled step by step from the papers above. Notation: Ψ\Psi parameters, LL layers, hidden size hh, aa heads, sequence length ss, microbatch size bb, global batch BB (sequences), degrees tt (TP), pp (PP), dd (DP), ee (EP), and m=B/(db)m = B/(db) microbatches per pipeline per step.

1. Parameters and model-state memory

A GPT-style layer has about 4h24h^2 attention weights (QQ, KK, VV, output) and 8h28h^2 MLP weights (h→4h→hh \to 4h \to h), so Ψ≈12Lh2\Psi \approx 12Lh^2, ignoring embeddings. TP and PP divide a rank’s share to Ψ/(tp)\Psi/(tp). Mixed-precision Adam stores 2 + 2 + 12 bytes per parameter, and ZeRO stage kk divides the first kk of {optimizer, gradients, weights} by dd:

Mstate=Ψtp[2d3+2d2+12d1]dk={dif stage≥k1otherwise\begin{gathered} M_{\mathrm{state}} = \frac{\Psi}{tp} \left[ \frac{2}{d_3} + \frac{2}{d_2} + \frac{12}{d_1} \right] \\ d_k = \begin{cases} d & \text{if stage} \ge k \\ 1 & \text{otherwise} \end{cases} \end{gathered}

For MoE, expert weights are further divided by ee and sharded over the d/ed/e replicas of each expert. The sim uses top-2 routing over half-width experts (4h24h^2 each), so active parameters per token equal the dense model’s.

2. Activation memory

Korthikanti et al. count the bytes stored per layer in 16-bit training. Attention needs 11sbh+5as2b11sbh + 5as^2b, the MLP 19sbh19sbh and the two LayerNorms 4sbh4sbh, for a total of

Alayer=sbh(34+5ash)A_{\mathrm{layer}} = sbh \left( 34 + \frac{5as}{h} \right)

With tensor and sequence parallelism every term divides by tt. Selective recomputation drops the 5as/h5as/h term (the softmax and attention-dropout tensors that grow with s2s^2). Full recomputation keeps only 2sbh2sbh per layer, or 2sbh/t2sbh/t if sharded across the TP group. In 1F1B the first stage holds pp microbatches of L/pL/p layers, so

Atotal≈sbhLt(34+5ash)A_{\mathrm{total}} \approx \frac{sbhL}{t} \left( 34 + \frac{5as}{h} \right)

This is independent of pp. For GPT-3 (a=96a = 96, s=2048s = 2048, h=12288h = 12288) the attention term is 5as/h=805as/h = 80, larger than 34, which is why selective recomputation pays off. The sim uses min⁡(p,m)⋅⌈L/p⌉\min(p, m) \cdot \lceil L/p \rceil layers, since a pipeline with fewer microbatches than stages never fills.

3. Ring collectives

A ring all-reduce over nn ranks of an SS-element buffer runs n−1n-1 reduce-scatter steps, in which each rank sends one S/nS/n chunk to its neighbor and adds the chunk it receives, then n−1n-1 all-gather steps that circulate the reduced chunks. Each element crosses 2(n−1)2(n-1) links in total, spread over nn links of bandwidth BB, so

tallreduce≈SB⋅2(n−1)ntRS=tAG≈SB⋅n−1n\begin{aligned} t_{\mathrm{allreduce}} &\approx \frac{S}{B} \cdot \frac{2(n-1)}{n} \\ t_{\mathrm{RS}} = t_{\mathrm{AG}} &\approx \frac{S}{B} \cdot \frac{n-1}{n} \end{aligned}

This is why a ring is bandwidth-optimal and nearly independent of nn, and also why it runs at the speed of its slowest link: put one scale-out hop in the ring and every step waits for it. Latency (≈2(n−1)\approx 2(n-1) message times) dominates for small SS, which is why TP’s many medium-sized all-reduces suffer most on a high-latency network.

chunk 1chunk 2chunk 3chunk 4rank 01642rank 14275rank 27531rank 33164sum: 15, 14, 20, 120123
1 / 7
Ranks n

Each of 4 ranks holds its own 4-chunk buffer. Goal: every rank ends with the element-wise sum.

Ring all-reduce: n − 1 reduce-scatter steps, then n − 1 all-gather steps. Every link carries one S/n chunk per step, so the time is ≈ (S/B) · 2(n − 1)/n whatever n is. Values illustrative.Share freely with credit: ‘Figure from chipfieldguide.com’

4. TP and PP volumes

Each TP all-reduce operates on a b×s×hb \times s \times h activation tensor, and there are four per layer per microbatch (two forward via gg, two backward via ff), so the volume per rank, in elements per layer per microbatch, is

VTP=4⋅2(t−1)t bsh=8bsh (t−1)tV_{\mathrm{TP}} = 4 \cdot \frac{2(t-1)}{t}\, bsh = \frac{8bsh\,(t-1)}{t}

Per step, multiply by L/pL/p and mm. With full recomputation the extra forward pass repeats its two all-reduces (×1.5). A PP stage boundary carries bshbsh elements per microbatch per direction, divided by tt with scatter/gather. Worked example (the sim’s default): s=2048s = 2048, h=12288h = 12288, t=8t = 8, L/p=12L/p = 12, m=64m = 64 gives 8×2048×12288×7/8×2 bytes≈352 MB8 \times 2048 \times 12288 \times 7/8 \times 2\,\mathrm{bytes} \approx 352\,\mathrm{MB} per layer per microbatch, about 270 GB per step. PP moves 2×64×2048×12288×2/8 bytes≈0.8 GB2 \times 64 \times 2048 \times 12288 \times 2/8\,\mathrm{bytes} \approx 0.8\,\mathrm{GB}, and DP, with 2.72×1092.72 \times 10^{9} local parameters and d=16d = 16, moves 2×5.4 GB×15/16≈10 GB2 \times 5.4\,\mathrm{GB} \times 15/16 \approx 10\,\mathrm{GB}.

5. The bubble

With flushes, the last stage starts p−1p-1 forward slots late and the first stage ends p−1p-1 backward slots early, so

tbubble=(p−1)(tf+tb)tideal=m (tf+tb)bubble fraction=p−1m\begin{aligned} t_{\mathrm{bubble}} &= (p-1)(t_{\mathrm{f}} + t_{\mathrm{b}}) \\ t_{\mathrm{ideal}} &= m\,(t_{\mathrm{f}} + t_{\mathrm{b}}) \\ \text{bubble fraction} &= \frac{p-1}{m} \end{aligned}

As a share of total step time it is (p−1)/(m+p−1)(p-1)/(m+p-1), GPipe’s O((K−1)/(M+K−1))O\big((K-1)/(M+K-1)\big). Interleaving vv chunks gives (p−1)/(vm)(p-1)/(vm). The sim reports (p−1)/m(p-1)/m in the Expert note and the share of step time in the status bar.

6. Compute time and exposed communication

A forward pass costs 24Bsh2+4Bs2h24Bsh^2 + 4Bs^2h FLOPs per layer, and the backward pass twice that, so without recomputation a step is about 72BsLh2 (1+s/(6h))72BsLh^2\,(1 + s/(6h)) FLOPs, or 96BsLh2 (1+s/(6h)+V/(16Lh))96BsLh^2\,(1 + s/(6h) + V/(16Lh)) with recomputation and the logit layer. The sim uses 6 ΨactiveBs+12LBs2h6\,\Psi_{\mathrm{active}} B s + 12 L B s^2 h (×4/3 for full recompute) divided by NN accelerators at an assumed 400 TFLOP/s. It then estimates

tstep≈tcompute(1+p−1m)+tTP+tEP+tPP+max⁡ ⁣(0,  tDP−23 tcompute)\begin{aligned} t_{\mathrm{step}} \approx{} & t_{\mathrm{compute}} \left( 1 + \frac{p-1}{m} \right) \\ & + t_{\mathrm{TP}} + t_{\mathrm{EP}} + t_{\mathrm{PP}} \\ & + \max\!\left( 0,\; t_{\mathrm{DP}} - \tfrac{2}{3}\, t_{\mathrm{compute}} \right) \end{aligned}

Here each tXt_X is volume ÷ bandwidth, at 450 GB/s if the group fits in a scale-up domain and 50 GB/s otherwise. The 23\tfrac{2}{3} reflects that the backward pass, which DP traffic can hide behind, is two-thirds of the compute. This is deliberately crude: real TP overlaps partly, real EP can be fully hidden (DualPipe), and latency is ignored.

7. A planning procedure

Narayanan et al.’s takeaways, turned into steps:

  1. Set tt to the smallest power of two up to the scale-up domain size that makes a layer’s state fit.
  2. Set pp so that the model-parallel size tptp fits weights, optimizer state and activations in memory, with ZeRO-1 or -2 on top.
  3. Let d=N/(tp)d = N/(tp) and check that m=B/(db)≫pm = B/(db) \gg p; if not, interleave or reduce pp.
  4. Tune bb for kernel efficiency against the bubble, and choose recomputation to fit memory at the largest useful bb.
  5. For MoE, size ee so experts fit and all-to-all stays mostly inside the domain or node-limited, as DeepSeek-V3 did.

Production systems search this space with measured cost models. Llama 3 notes it built a memory estimator and a performance-projection tool to explore configurations.

Novice · 0 of 4 correct
  1. Q1Training with the Adam optimizer in mixed precision takes about 16 bytes per parameter for weights, gradients and optimizer state. Roughly how much memory is that for a 7.5-billion-parameter model, before activations?

  2. Q2Which parallelism strategy sends the least data per step for a typical large model, and why?

  3. Q3With 8 pipeline stages and 8 microbatches per step, about what share of each stage’s time is idle bubble?

  4. Q4Llama 3 405B was trained with the parallelism order [TP, CP, PP, DP], innermost first. What does ‘innermost’ mean for the hardware?

Sources

Show Hide 16 sources
  1. The Llama 3 Herd of ModelsLlama Team, AI @ Meta · arXiv:2407.21783 · 2024405B model trained on up to 16K H100 80 GB GPUs, eight per server joined by NVLink; 4D parallelism ordered [TP, CP, PP, DP] with the innermost kept inside a server; Table 4 configurations and 38–43% MFU; 24K-GPU RoCE cluster, 3,072-GPU pods at full bisection, 1:7 oversubscription above; 466 interruptions in 54 days.
  2. ZeRO: Memory Optimizations Toward Training Trillion Parameter ModelsSamyam Rajbhandari, Jeff Rasley, Olatunji Ruwase, Yuxiong He · arXiv:1910.02054 · 2019Mixed-precision Adam needs 2Ψ + 2Ψ + 12Ψ = 16Ψ bytes; GPT-2 1.5B needs at least 24 GB versus 3 GB of fp16 weights; three partitioning stages; communication 2Ψ for standard DP and stages 1–2, 3Ψ (1.5×) with parameter partitioning; 40B model with Megatron across two DGX-2 nodes ran at about 5 TFLOPS per V100.
  3. Megatron-LM: Training Multi-Billion Parameter Language Models Using Model ParallelismMohammad Shoeybi, Mostofa Patwary, Raul Puri, Patrick LeGresley, Jared Casper, Bryan Catanzaro · arXiv:1909.08053 · 2019Column- then row-parallel split of the MLP and attention; f and g operators; two all-reduces in the forward pass and two in the backward pass per layer; 8.3B model on 512 V100s with 8-way model parallelism at 15.1 PFLOPS and 76% scaling efficiency; 300 GB/s inside a DGX-2H versus 100 GB/s between servers.
  4. Efficient Large-Scale Language Model Training on GPU Clusters Using Megatron-LMDeepak Narayanan, Mohammad Shoeybi, Jared Casper, et al. · arXiv:2104.04473 (SC21) · 2021PTD-P; bubble fraction (p−1)/m and (1/v)(p−1)/m interleaved; 1F1B stashes at most p microbatches; TP moves 8bsh(t−1)/t per layer per microbatch, PP moves bsh; Takeaways 1–3; 1T model on 3072 A100s (t = 8, p = 64) at 502 PFLOPS, 52% of peak; 8×200 Gb/s HDR InfiniBand per node; scatter/gather; PP 892 GB/s vs DP 12.9 TB/s effective bisection bandwidth; 70% faster than ZeRO-3.
  5. Reducing Activation Recomputation in Large Transformer ModelsVijay Korthikanti, Jared Casper, Sangkug Lym, et al. · arXiv:2205.05198 · 2022Activation memory per layer sbh(34 + 5as/h) bytes; divided by t with tensor plus sequence parallelism; first pipeline stage stores L layers’ worth; Table 2 for recompute options; sequence parallelism swaps all-reduce for all-gather plus reduce-scatter at the same bandwidth; 5as/h = 80 for GPT-3.
  6. GPipe: Easy Scaling with Micro-Batch Pipeline ParallelismYanping Huang, Youlong Cheng, Ankur Bapna, et al. · arXiv:1811.06965 · 2018Micro-batch pipelining; bubble O((K−1)/(M+K−1)); negligible when M ≥ 4K; communication only at partition boundaries.
  7. GShard: Scaling Giant Models with Conditional Computation and Automatic ShardingDmitry Lepikhin, HyoukJoong Lee, Yuanzhong Xu, et al. · arXiv:2006.16668 · 2020Mixture-of-experts Transformer beyond 600B parameters on 2048 TPU v3 in 4 days; MoE dispatch and combine are AllToAll; AllToAll cost about O(√D) on a 2D TPU mesh; going from 128 to 2048 experts raised its share of time from 16% to 36%.
  8. Switch Transformers: Scaling to Trillion Parameter Models with Simple and Efficient SparsityWilliam Fedus, Barret Zoph, Noam Shazeer · arXiv:2101.03961 (JMLR 2022) · 2021Sparse MoE keeps compute per token constant as parameters grow; expert parallelism allocates the data-parallel cores to experts; tokens move by all-to-all; combining expert and model parallelism adds all-reduce on top.
  9. DeepSeek-V3 Technical ReportDeepSeek-AI · arXiv:2412.19437 · 2024671B total, 37B active parameters; 2048 H800 GPUs, 8 per node on NVLink, InfiniBand between nodes; 16-way PP, 64-way EP over 8 nodes, ZeRO-1 DP, no TP; NVLink 160 GB/s vs IB 50 GB/s; tokens limited to 4 nodes; compute-to-communication about 1:1, DualPipe overlap; 2.788M GPU hours.
  10. NCCL tests: Performance reported by NCCL tests (PERFORMANCE.md)NVIDIA · NVIDIA/nccl-tests (GitHub)Ring all-reduce needs 2(n−1) transfers per element; t = (S/B)·2(n−1)/n; reduce-scatter and all-gather each (n−1)/n; definition of bus bandwidth.
  11. Horovod: fast and easy distributed deep learning in TensorFlowAlexander Sergeev, Mike Del Balso · arXiv:1802.05799 · 2018Ring-allreduce: each of N nodes talks to two peers 2(N−1) times, reducing then distributing chunks; bandwidth-optimal per Patarasuk and Yuan; no parameter server.
  12. PyTorch Distributed: Experiences on Accelerating Data Parallel TrainingShen Li, Yanli Zhao, Rohan Varma, et al. · arXiv:2006.15704 (VLDB 2020) · 2020Gradient bucketing; all-reduce launched during the backward pass to overlap communication with computation; skipping synchronization; near-linear scaling on 256 GPUs.
  13. PyTorch FSDP: Experiences on Scaling Fully Sharded Data ParallelYanli Zhao, Andrew Gu, Rohan Varma, et al. · arXiv:2304.11277 (VLDB 2023) · 2023Parameters gathered on demand before computation and discarded afterwards; forward and backward prefetching of AllGather; gradient accumulation with or without communication.
  14. Rail-only: A Low-Cost High-Performance Network for Training LLMs with Trillion ParametersWeiyang Wang, Manya Ghobadi, Kayvon Shakeri, Ying Zhang, Naader Hasani · arXiv:2307.12169 · 2023High-bandwidth (HB) domains of 8 (DGX H100, MI300X), 72 (GB200 NVL72) and 256 (DGX GH200) GPUs; 400 Gb/s NIC per GPU; GPT-1T traffic ~300 GB per GPU pair inside an HB domain vs ~6 GB in the NIC domain; traffic leaving HB domains stays within rails; 38–77% lower network cost; 8.2–11.2% overhead for MoE all-to-all.
  15. TPU v4: An Optically Reconfigurable Supercomputer for Machine Learning with Hardware Support for EmbeddingsNorman P. Jouppi, George Kurian, Sheng Li, et al. · arXiv:2304.01433 (ISCA 2023) · 20234096-chip supercomputer; optical circuit switches; topology chosen per job to match parallelism (e.g. 4×4×32 for pipeline); data parallelism along one torus dimension; 1.2–2.3× gains from topology and hyperparameters; six ICI links at 50 GB/s.
  16. Ring Attention with Blockwise Transformers for Near-Infinite ContextHao Liu, Matei Zaharia, Pieter Abbeel · arXiv:2310.01889 · 2023Splits long sequences across devices in a ring and overlaps passing key-value blocks with blockwise attention, with no added overhead while compute takes longer than transfer.