← All Posts
Deep Learning · Popular Videos · Umar Jamil· Part 14 · Chapters 27–30

Expert Parallelism & All-to-All

Sparse activation does not make the parameters go away. A 671B-parameter MoE still needs 671B parameters stored somewhere, even if each token touches 37B of them. Expert parallelism puts different experts on different devices — and since tokens must then travel to their experts, introduces the one collective this series has not yet used: all-to-all.

Why MoE forces a new axis

Data parallelism replicates. For a dense model that is merely wasteful; for an MoE it is impossible, because the total parameter count is the thing that grew. The other axes do not fit either:

AxisApplied to an MoE layer
Data parallelReplicates all $E$ experts on every device. Defeats the point entirely.
Tensor parallelWorks — shard each expert's matrices — but every device still stores a slice of every expert, and the shards become small and inefficient.
Pipeline parallelSplits by layer, so an individual MoE layer is still whole on one device.
Expert parallelGive each device $E/P$ whole experts. Memory divides by $P$ and each expert matmul stays full size.
Expert parallelism shards the layer by expert index. Non-expert weights — attention, norms, the router — stay replicated as usual. Only the expert bank is distributed.

The catch is immediate. A token arriving on device 0 may be routed to an expert living on device 5. Either the expert comes to the token or the token goes to the expert, and experts are far larger than tokens.

All-to-all: the distributed transpose

The routing produces, on every device, a pile of tokens destined for every other device. Every rank has data for every rank — which is precisely the all-to-all pattern.

Dispatchall-to-all: send each token to the rank owning its expert
→
Computeeach rank runs its own experts on whatever arrived
→
Combineall-to-all: send the outputs back to the originating rank

The combine is the exact inverse permutation of the dispatch. Once results return, each token's $k$ expert outputs are weighted by their gate values and summed into the residual stream.

The backward pass needs no new machinery. All-to-all is a permutation of data, and the adjoint of a permutation is its inverse. From the duality table, the backward of an all-to-all is another all-to-all with source and destination swapped — so backward dispatch is forward combine, and vice versa.

Why fixed-size buffers matter here

An all-to-all is far more efficient when every rank sends the same number of bytes to every other rank, because the transfer sizes are known in advance and no metadata exchange is needed. Routing, however, is data-dependent and uneven.

This is exactly what the capacity factor from the previous note buys: it fixes the per-expert buffer at $C=\text{cf}\cdot kT/E$ tokens, so the all-to-all becomes a uniform, statically shaped collective. The price is dropped tokens on overflow and padded compute on underflow.

Load balance is now a systems problem, not only a quality problem. The all-to-all is synchronous: every rank waits for the slowest. A skewed router makes one device the bottleneck for the entire group, so balancing losses pay for themselves in throughput long before they pay for themselves in loss.

The bandwidth bill

Each rank sends a $(P-1)/P$ fraction of its routed tokens, each a vector of width $d$. With $T$ tokens per rank and top-$k$ routing, one dispatch moves

$$kT\,d\cdot\frac{P-1}{P}\cdot\text{bytes}.$$

Dispatch and combine double it, and the backward pass doubles it again — four all-to-alls per MoE layer per step. For $d=7168$, $T=8192$, $k=8$, BF16, across 58 layers:

EP degreeOne dispatchPer layer per stepWhole model per step
8784 MiB3.1 GiB178 GiB
16840 MiB3.3 GiB190 GiB
64882 MiB3.4 GiB200 GiB
256892 MiB3.5 GiB202 GiB
The volume barely grows with $P$. The factor $(P-1)/P$ saturates at 1, exactly as it did for all-reduce. Widening the expert-parallel group is nearly free in bytes — what changes is which links those bytes cross, and that is what actually decides the cost.

Hence node-limited routing: constrain each token to experts residing on at most a small number of nodes, so most all-to-all traffic stays on the fast intra-node fabric. It is a routing constraint adopted for a purely topological reason.

experts per device
expert memory / device
all-to-all / layer / step
tokens leaving each rank
256 experts of width 2048 over a $d=7168$ model, 8192 tokens per rank, top-8 routing, BF16. Colours identify the rank a token originated on, so you can watch it leave and return.

Tensor parallelism for experts

Expert parallelism divides memory by $P$ but cannot exceed $E$, and with few experts each device may still hold too much. Tensor parallelism applies to an expert exactly as it does to a dense FFN — the gate and up projections column-parallel, the down projection row-parallel, one all-reduce at the end.

EP: split by expert

Whole experts on separate devices. Matmuls stay large and efficient. Communication is all-to-all on tokens. Limited by $E$.

TP: split inside an expert

Every device holds a slice of every expert. Communication is all-reduce on activations. Matmuls shrink and can become inefficient.

Expert tensor parallelism

The two are orthogonal, so combine them: arrange $P=P_{\text{EP}}\times P_{\text{TP}}$ as a two-dimensional grid. Each EP group owns a distinct set of experts, and within a group the TP ranks shard those experts' matrices.

$$\text{expert memory per device}=\frac{E}{P_{\text{EP}}}\cdot\frac{3\,d\,d_{\text{ff}}}{P_{\text{TP}}}.$$

The communication pattern nests: an all-to-all across the EP axis to place tokens, then an all-reduce within the TP axis to finish each expert's matmul.

SymptomAdjust
Expert matmuls too small to reach peakLower $P_{\text{TP}}$, raise $P_{\text{EP}}$
$P_{\text{EP}}$ already equals $E$ and memory still tightRaise $P_{\text{TP}}$ — the only remaining option
All-to-all crossing node boundariesKeep the EP group within a node, or use node-limited routing
One rank consistently lateFix routing balance before touching the parallelism layout
Takeaway

Expert parallelism shards the expert bank and pays for it with two all-to-alls per layer, whose volume barely grows with group size but whose cost depends entirely on which links it crosses. Fixed capacity buffers make the collective statically shaped, and tensor parallelism composes on a second axis when the expert count runs out.

Check yourself