Expert Parallelism & 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:
| Axis | Applied to an MoE layer |
|---|---|
| Data parallel | Replicates all $E$ experts on every device. Defeats the point entirely. |
| Tensor parallel | Works — shard each expert's matrices — but every device still stores a slice of every expert, and the shards become small and inefficient. |
| Pipeline parallel | Splits by layer, so an individual MoE layer is still whole on one device. |
| Expert parallel | Give each device $E/P$ whole experts. Memory divides by $P$ and each expert matmul stays full size. |
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.
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.
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.
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
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 degree | One dispatch | Per layer per step | Whole model per step |
|---|---|---|---|
| 8 | 784 MiB | 3.1 GiB | 178 GiB |
| 16 | 840 MiB | 3.3 GiB | 190 GiB |
| 64 | 882 MiB | 3.4 GiB | 200 GiB |
| 256 | 892 MiB | 3.5 GiB | 202 GiB |
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.
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.
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.
| Symptom | Adjust |
|---|---|
| Expert matmuls too small to reach peak | Lower $P_{\text{TP}}$, raise $P_{\text{EP}}$ |
| $P_{\text{EP}}$ already equals $E$ and memory still tight | Raise $P_{\text{TP}}$ — the only remaining option |
| All-to-all crossing node boundaries | Keep the EP group within a node, or use node-limited routing |
| One rank consistently late | Fix routing balance before touching the parallelism layout |
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
- Explain why data parallelism cannot be used for the expert bank. reasoning
- Derive the backward collective for an all-to-all from the duality table. derivation
- Compute the all-to-all volume per layer for $d=4096$, $T=4096$, $k=2$, $P=16$. calculation
- Explain why the capacity factor is a systems requirement here, not just a quality knob. reasoning
- Given $E=64$ experts and 128 devices, propose an EP/TP split and justify it. design