skip to content

How does expert parallelism shard an MoE model for serving across GPUs?

level: seniorimportance: should knowfreq 38%

answer

  1. whole experts, not matrix slices
  2. the collective changes shape
  3. every rank talks to every rank
  4. the busiest expert sets the pace
  5. attention and experts sharded differently

basics

~20 s

Expert parallelism places whole experts on different GPUs instead of slicing every matrix. Each MoE layer then does an all-to-all: tokens are dispatched to the GPUs owning their chosen experts and the results are gathered back. Uneven routing makes one GPU the straggler.

solid answer

~60 s

For a mixture-of-experts checkpoint, most of the parameter mass sits in expert weights that any given token barely touches. Expert parallelism exploits that by assigning whole experts to GPUs — GPU 0 holds experts 0-15, GPU 1 holds 16-31 — so memory divides cleanly without slicing each matrix thin. The communication pattern changes accordingly: instead of an all-reduce, each MoE layer performs an **all-to-all** to dispatch tokens to whichever ranks own their routed experts, and a second all-to-all to combine the outputs. Two consequences dominate in production. First, load imbalance: the router is not uniform, so a hot expert makes its GPU the straggler and every other rank waits on it at the layer barrier — throughput tracks the busiest rank, not the average. Second, layouts are usually hybrid, with attention run tensor-parallel while the MoE blocks run expert-parallel, since the attention weights are small and shard well while the experts do not. Engines expose it separately from the tensor-parallel degree — vLLM, for example, has an `--enable-expert-parallel` switch alongside `--tensor-parallel-size`.

go deeper

for a junior

Know that some models are built from many experts and that serving them can place different experts on different GPUs rather than splitting every weight matrix.

for a middle

Explain the mechanics: whole experts per rank, an all-to-all to dispatch tokens and another to combine results, and why that differs from an all-reduce.

for a senior

Own the operational reality — diagnose hot-expert imbalance from per-expert token counts and per-rank utilization, and fix it with replication, placement or larger batches.

for a principal

Decide the hybrid layout for the fleet: which parts run tensor-parallel and which expert-parallel, what fabric that demands, and whether the memory saving justifies the all-to-all traffic and the imbalance risk.

## The shape of the problem A mixture-of-experts checkpoint has a large total parameter count but activates only a fraction of it per token. That is wonderful for compute and awkward for memory: you still have to hold every expert somewhere, because any token may route to any of them. Serving such a model across GPUs therefore has a natural split that dense models do not: distribute *experts*, not matrix slices. ## What expert parallelism does Each GPU is made the owner of a subset of the experts in every MoE layer. Attention and other dense parts are handled by whatever strategy you chose for them. When a token reaches an MoE layer, its router picks the experts it should visit; if those experts live on other GPUs, the token's hidden state has to travel there, be processed, and come back. That gives the pattern its signature collective: an **all-to-all**. Every rank sends each other rank the tokens destined for the experts it owns (dispatch), the experts run locally, and a second all-to-all returns the outputs to the ranks that owned the tokens (combine). Two all-to-alls per MoE layer, per forward pass. ## Why it is preferred over slicing each expert You *could* apply tensor parallelism to expert weights, cutting every expert matrix across all GPUs. It works, but it produces very thin per-rank slices — small, inefficient matmuls with poor kernel utilization — and it makes every rank participate in every expert. Expert parallelism keeps each expert's matmul whole and reasonably sized, which is friendlier to the hardware, and it splits the bulk of the parameters with no per-matrix collective at all. The usual production layout is hybrid: tensor-parallel attention (small weights, latency-critical, shards well) combined with expert-parallel MoE blocks (large weights, naturally partitionable). ## Load imbalance is the operational story Routing is learned, not uniform. In any real traffic mix some experts are chosen far more often than others, and because the layer ends in a barrier, every rank waits for the slowest. A single hot expert therefore caps the throughput of the entire group while other GPUs idle. This is the failure mode interviewers are probing for, and the mitigations are worth naming: - **Placement.** Spread experts that co-activate, rather than assigning contiguous blocks blindly. - **Replication of hot experts.** Duplicate the busiest experts across several ranks so their traffic splits, spending memory to buy balance. - **Larger batches.** More tokens per step average out routing noise; at small batch the imbalance is extreme. - **Measurement first.** Per-expert token counts under representative traffic tell you whether you have an imbalance problem at all; without them you are guessing. ## Communication characteristics All-to-all is more sensitive to topology than all-reduce, because every rank talks to every other rank rather than following an optimized ring or tree. Its volume also scales with the number of tokens in flight and with how many experts each token visits. On a well-connected node this is manageable; across a slow link it is punishing, which is why expert-parallel groups are normally kept inside one high-bandwidth domain, exactly as tensor-parallel groups are. Larger cross-node expert-parallel deployments exist for very large MoE models, but they demand a serious fabric and careful placement. ## Sizing intuition Because each rank holds only its experts, memory per GPU falls roughly with the expert-parallel degree for the expert weights — the dominant term — while attention weights and the KV cache follow whatever strategy you applied to them. Note that the KV cache is *not* affected by expert parallelism at all: it belongs to attention. If cache capacity is your constraint, expert parallelism does not solve it; tensor-parallel attention or a different cache dtype does. ## Where the engines put it The degree and the switch are engine-specific. vLLM exposes `--enable-expert-parallel` as a distinct option from `--tensor-parallel-size`, so you choose the MoE strategy independently of the dense one. Other stacks fix the layout when the engine is built. Do not assume a flag name carries between engines. ## Weak answers "It is just tensor parallelism for MoE models" misses the different collective and the imbalance problem entirely. "It reduces compute per GPU" confuses sparsity — a property of the model — with sharding, which redistributes memory and adds communication.

  • Why is all-to-all harder on the interconnect than all-reduce?
    All-reduce has a fixed, symmetric communication pattern that libraries optimize into rings or trees with predictable per-rank volume. All-to-all has every rank sending a different, data-dependent amount to every other rank, so it cannot be pipelined as neatly and is far more sensitive to topology and to imbalance in the message sizes. Skewed routing turns some of those messages into stragglers that delay the whole layer.
  • Does expert parallelism reduce the KV cache each GPU must hold?
    No. The KV cache belongs to the attention layers, and expert parallelism only redistributes the feed-forward experts. If cache capacity is what limits your concurrency or context length, you need tensor-parallel attention, a smaller cache dtype, or fewer concurrent sequences. This is a common confusion because expert parallelism dramatically reduces *weight* memory per GPU, which makes the remaining cache pressure more visible, not less.
  • How would you detect that expert load imbalance is costing you throughput?
    Instrument per-expert token counts over representative traffic and compare the busiest expert's share against the uniform expectation; a large skew is the signal. Corroborate with per-rank GPU utilization — imbalance shows up as one device pinned while its peers idle at the layer barrier. Then act by replicating the hot experts across ranks, revisiting placement, or increasing batch size so routing averages out.

saying these in an interview costs you the question

  • Describing expert parallelism as tensor parallelism applied to experts
  • Assuming router traffic is evenly spread across experts
  • Claiming it shrinks the KV cache per GPU
  • Expecting all-to-all to behave like an all-reduce on a slow link
  • Setting the expert-parallel layout without measuring per-expert load

context