Mixture of Experts as a model is covered in Section XIII.02. This file is about the distributed systems problem it creates.
1. What is it?#
In a Mixture-of-Experts model, each layer has E expert FFNs, and each token is routed to k of them (typically k=1 or 2). Expert parallelism places different experts on different GPUs.
GPU 0 GPU 1 GPU 2 GPU 3
┌──────────┐ ┌──────────┐ ┌──────────┐ ┌──────────┐
│ experts │ │ experts │ │ experts │ │ experts │
│ 0-15 │ │ 16-31 │ │ 32-47 │ │ 48-63 │
└──────────┘ └──────────┘ └──────────┘ └──────────┘
▲ ▼ ▲ ▼ ▲ ▼ ▲ ▼
└───────── ALL-TO-ALL ─────────────────────────────┘
(send each token to the GPU holding its chosen experts,
then send the results back)2. Problem → Why → Optimization#
PROBLEM An MoE model has enormous total parameters (e.g. DeepSeek-V3: 671B
total, 37B active per token). The weights don't fit on one GPU, but
replicating all experts everywhere wastes memory.
WHY Only k of E experts are used per token, so most experts are idle
for any given token — but you don't know which until the router runs.
OPTIMIZE Distribute experts across GPUs. Route tokens to the GPU holding
their expert, compute, route the results back.
TRADE-OFFS
✓ memory scales: E experts across N GPUs = E/N per GPU
✓ compute is sparse: only k experts run per token
✗ ALL-TO-ALL communication, twice per MoE layer
✗ LOAD IMBALANCE: routing is data-dependent and uneven
✗ the All-to-All is a synchronization point — the slowest GPU sets the pace
WHEN TO USE Serving MoE models that don't fit with TP alone.
WHEN NOT TO Dense models (no experts), or when the model fits with TP.3. Simple analogy#
A hospital with specialists in different buildings.
Each patient (token) needs to see 2 of 64 specialists. The specialists are distributed across 4 buildings. So:
- Look at each patient’s needs (the router).
- Send every patient to the building(s) holding their specialists (All-to-All).
- Specialists see their patients.
- Send every patient back to where they started (All-to-All again).
The problems are obvious in the analogy: the transport is expensive, and if 60% of patients need a specialist in building 2, that building is overwhelmed while others idle. Both problems are real and both are why EP is harder than TP.
4. The All-to-All#
Before dispatch: each GPU has a batch of tokens, in arbitrary order.
Router output: for each token, which expert(s) it needs.
DISPATCH (All-to-All #1):
1. Sort tokens by destination expert (and therefore by destination GPU)
2. Count how many tokens go to each GPU
3. All-to-All: exchange counts, then exchange the token data
4. Each GPU now has a contiguous set of tokens for each of ITS experts
COMPUTE:
Grouped GEMM: each expert processes its tokens.
(Section IV.03 — many different-shaped GEMMs in one launch)
COMBINE (All-to-All #2):
Send the results back to the originating GPUs, in the original order.
Weighted-sum the k expert outputs per token.Note step 1 in dispatch: sorting tokens by expert. This is essential — it turns what would be maximally divergent per-token work (Section VI.10) into a grouped GEMM with contiguous memory access. Any MoE implementation that doesn’t sort is leaving a large factor on the table.
The communication volume#
Per MoE layer, per token, with top-k routing:
dispatch: k × d × bytes sent
combine: k × d × bytes returned
Mixtral-8x7B: d=4096, k=2, BF16, 32 MoE layers, batch 64:
per layer: 64 tokens × 2 experts × 4096 × 2 bytes × 2 (dispatch+combine)
= 2.1 MB
× 32 layers = 67 MB per decode step
Over NVLink (450 GB/s effective): 0.15 ms — fine
Over InfiniBand 400G (50 GB/s): 1.3 ms — significant
Over PCIe (25 GB/s): 2.7 ms — problematicAll-to-All is worse than AllReduce for a given volume because it has no ring structure to exploit — every GPU sends distinct data to every other GPU. Its cost scales with N in a way AllReduce’s doesn’t.
5. Load imbalance — the hard problem#
Routing is learned and data-dependent. Some experts are more popular:
Ideal (perfectly balanced), 8 experts, 512 tokens:
each expert gets 64 tokens
Reality:
expert 0: 140 tokens ← "hot" expert
expert 1: 95
expert 2: 78
expert 3: 62
expert 4: 48
expert 5: 40
expert 6: 31
expert 7: 18 ← "cold" expert
The GPU holding expert 0 takes 140/64 = 2.2x the time.
ALL other GPUs wait for it at the All-to-All barrier.
→ effective utilization = 64/140 = 46%The All-to-All is a barrier, so the slowest GPU determines the step time. This is the dominant inefficiency in MoE serving.
Mitigations#
1. CAPACITY FACTOR
Cap tokens per expert at `capacity = C × tokens/E`. Excess tokens are
DROPPED (skip the expert, pass through the residual).
C=1.25 is typical. Bounds the imbalance; costs a little quality.
2. EXPERT REPLICATION
Replicate hot experts on multiple GPUs. Requires knowing which are hot
(measurable offline) and adds memory.
3. DYNAMIC EXPERT PLACEMENT
Rebalance expert-to-GPU assignment based on measured load.
[EMERGING] — DeepSeek's inference system does a version of this.
4. LARGER BATCHES
Imbalance averages out with more tokens. At batch 4096 the distribution
is much closer to uniform than at batch 64.
→ EP strongly favors high-batch serving.
5. AUXILIARY LOSS DURING TRAINING
Encourages balanced routing. Helps but doesn't eliminate imbalance.Point 4 is the practical one for serving: MoE models are much more efficient at large batch, because that’s when the routing distribution smooths out. At batch 1, one expert per token means you’re using k/E of your FFN weights and paying full All-to-All latency. MoE is a throughput-oriented architecture.
6. Combining EP with TP#
Real deployments combine them:
DeepSeek-V3-scale model, 8 nodes × 8 GPUs = 64 GPUs:
Attention: TP=8 within each node (attention is dense)
Experts: EP=64 across all GPUs (each GPU holds 256/64 = 4 experts)
Per layer:
- attention: 1 AllReduce within the node (NVLink, fast)
- MoE: All-to-All across all 64 GPUs (crosses nodes — expensive)The All-to-All crossing node boundaries is the bottleneck in large MoE deployments, and it’s why DeepSeek and others invest heavily in overlapping communication with computation (running the attention of the next layer while the MoE All-to-All is in flight).
Alternative: expert TP#
For smaller MoE models, you can shard each expert with TP instead of distributing whole experts:
EP: expert 0 entirely on GPU 0 → All-to-All, load imbalance
ETP: expert 0 sharded across GPUs → AllReduce, no imbalance, but all
experts' weights on every GPU (no
memory saving)Use ETP when the model fits and you want to avoid All-to-All; use EP when memory forces it.
7. Performance#
Mixtral-8x7B (47B total, 13B active), 4×A100:
Configuration Memory/GPU ITL (b=1) Throughput (b=64)
TP=4 23.5 GB 12 ms 1,850 tok/s
EP=4 23.5 GB 18 ms 1,620 tok/s
EP=4, batch 512 23.5 GB — 4,900 tok/s ← EP shines at high batch
Note: for a model this size, TP is fine and simpler. EP matters for models
where E is large and the experts genuinely don't fit.For Mixtral-scale models, TP is usually the better choice. EP becomes necessary at DeepSeek-V3 scale (256 experts, 671B total).
8. Production implications#
- Use TP for MoE models that fit. EP only when memory forces it.
- EP favors high batch. If your traffic is low-concurrency, MoE + EP is a poor fit.
- Keep the All-to-All within a node if at all possible. Cross-node All-to-All is the dominant cost in large MoE deployments.
- Monitor expert load distribution. It’s a first-class metric for MoE serving; a skewed distribution directly costs throughput.
- Set the capacity factor deliberately. Too low drops tokens (quality); too high wastes compute.
- Grouped GEMM kernel quality matters a lot. A naive per-expert loop is far slower than a proper grouped GEMM.
- MoE models have a memory/compute asymmetry: cheap per token (few active params), expensive in memory (all params resident). They suit high-throughput serving on memory-rich hardware.
9. Common mistakes#
Using EP when TP would fit. Unnecessary complexity and All-to-All cost.
Not sorting tokens by expert. Maximal warp divergence and uncoalesced access.
Ignoring load imbalance. It can halve your throughput invisibly.
Cross-node All-to-All without overlap. Dominates the step time.
EP at low batch. All the cost, none of the amortization.
Capacity factor too low. Silently drops tokens and degrades quality.
Assuming MoE’s “13B active params” means it serves like a 13B model. It has 47B of memory traffic for weight residency and All-to-All overhead. It serves like something in between.
10. Hands-on exercise#
A. Measure the imbalance. Run an MoE model and instrument the router to record tokens per expert. Plot the distribution at batch 1, 16, 128, 1024. How does the imbalance change with batch size?
B. Compute the All-to-All cost. For an MoE model you can run, compute the per-token
All-to-All volume. Measure the actual All-to-All time with nccl-tests (alltoall_perf) for
that message size. What fraction of the step is it?
C. Sorting matters. Implement expert dispatch (i) with a naive per-token loop and (ii) with sort-then-grouped-GEMM. Measure both.
D. Capacity factor. Sweep the capacity factor from 1.0 to 2.0. Measure throughput and the token-drop rate. Where’s the knee?
E. TP vs EP. For an MoE model that fits either way, compare TP and EP on throughput at batch 16 and batch 512. Confirm EP’s advantage grows with batch.
11. Interview questions#
- What is expert parallelism and what collective does it use?
- Why is All-to-All harder to make fast than AllReduce?
- Explain MoE load imbalance and its effect on throughput.
- Name four mitigations for load imbalance.
- Why does EP favor high-batch serving?
- Why must tokens be sorted by expert before the GEMM?
- When would you choose TP over EP for an MoE model?
- Does a “671B total, 37B active” model serve like a 37B model? Explain.
12. Further reading#
- [ESTABLISHED] Lepikhin et al., “GShard” (2020); Fedus et al., “Switch Transformers” (2021)
- [ESTABLISHED] DeepSeek-V3 technical report — large-scale MoE inference engineering
- [ESTABLISHED] Rajbhandari et al., “DeepSpeed-MoE” (2022)
- [REFERENCE] NCCL All-to-All; CUTLASS grouped GEMM
- Next: 06 — Sequence and context parallelism