Serving a MoE across GPUs
Experts are spread over many GPUs, so every token's data has to be shipped to the machines holding its experts and shipped back, which makes the network and the busiest expert the real bottlenecks.
- 12 min read
- 3 reading levels
- Updated
Read these first
On this page 9
One lesson, three depths. Pick the one that fits you today — you can switch any time.
Beginner — No maths. Plain English.
The short answer
The experts do not fit on one machine, so every word travels to the machines holding its experts and back.
The food court
You are at a food court with thirty-two counters. Your order has items from eight of them.
Somebody has to walk to eight counters, wait at each, collect the items, and bring everything back to your table. You cannot start eating until the last item arrives.
Two things go wrong. All that walking takes time. And if one counter is popular, everyone is stuck in that queue while the other counters stand idle.
Serving a mixture-of-experts model is that food court, thousands of times a second.
Why the experts are spread out
A large sparse model can need over a thousand gigabytes to hold. No single card comes close.
So the experts are dealt out across many cards, a few each. Each card holds its share and nothing else.
Then every word must reach the cards holding its chosen experts. That trip, out and back, happens at every layer.
token arrives on GPU 3
|
+--> send its data to GPU 7, GPU 12, GPU 19, ... (its experts)
|
| each of those runs its expert on it
|
+<-- send the results back to GPU 3
|
combine, continue to the next layerThe two problems, in order of pain
Waiting for the slowest counter. All the cards must finish before the layer is done. One overloaded card holds up all the others. In the measurement below, a realistic traffic pattern leaves nearly half the machine idle.
The walking itself. Data goes out and comes back, at every layer, for every word. It adds up to a lot of network traffic per batch.
What is done about it
Open a second counter for the popular item. Keep a spare copy of the busiest experts on other cards and split their traffic. This is the biggest single win, and it needs measurements of who is actually busy.
Do not let an order span the whole court. Restrict each word to a few nearby machines before choosing experts within them. Fewer trips, shorter distances, very little quality lost.
Send the next order while cooking this one. Overlap the travel with the work so the network time hides behind the arithmetic.
The part that surprises people
The model's clever design is not the hard part of running one. The hard part is the traffic.
At this size, an architecture decision is a networking decision. That is why recent models restrict routing to a handful of machines: not for accuracy, but for the cables.
Where you have already seen this
- Hosted models that are cheap per word but need a whole cluster.
- Serving frameworks with settings named for expert parallelism and load balancing.
- Open models that are easy to download and hard to actually run.
Remember this
- Experts live on different machines. Every word travels to its experts and back, every layer.
- The busiest machine sets the pace, so uneven traffic wastes the whole cluster.
- The fixes are copying hot experts, restricting how far a word may travel, and overlapping travel with work.
What to learn next
- A working MoE layer in PyTorch — the layer these systems are distributing.
- Load balancing and expert collapse — the training-time half of the same problem.
- Model serving — the surrounding infrastructure this plugs into.
Developer — Code and libraries.
Setup
pip install torchWritten against PyTorch 2.5.1, Python 3.10. CPU, a few seconds. This simulates the placement problem; it does not need a cluster.
Where the time actually goes
import torch
torch.manual_seed(0)
T, E, K, DEV = 8192, 256, 8, 32 # tokens, experts, top-k, GPUs
h, nbytes = 7168, 2 # DeepSeek-V3 hidden size, bf16
# Skewed real-world routing: a few experts are genuinely popular.
pref = torch.softmax(torch.randn(E) * 0.8, dim=0)
choice = torch.multinomial(pref.expand(T, E), K, replacement=False)
load = torch.bincount(choice.flatten(), minlength=E).float()
print(f"{E} experts on {DEV} GPUs, {E//DEV} experts each, {T:,} tokens, top-{K}")
def device_loads(placement):
"""placement: expert -> list of devices holding a copy; load is split evenly."""
per_dev = torch.zeros(DEV)
for e in range(E):
devs = placement[e]
for d in devs:
per_dev[d] += load[e] / len(devs)
return per_dev
plain = {e: [e % DEV] for e in range(E)} # round-robin, one copy each
d1 = device_loads(plain)
print(f"\nround-robin placement, no replication")
print(f" busiest GPU {d1.max():8.0f} tokens | mean {d1.mean():8.0f} | "
f"max/mean {d1.max()/d1.mean():.2f}")
print(f" every GPU waits for the slowest -> {1 - d1.mean()/d1.max():.0%} of GPU time is idle")
# EPLB-style: replicate the hottest experts onto spare slots.
R = 32 # 32 redundant expert slots
hot = load.argsort(descending=True)[:R]
repl = {e: [e % DEV] for e in range(E)}
for i, e in enumerate(hot.tolist()):
repl[e] = repl[e] + [(e + 1 + i) % DEV] # a second home for each hot expert
d2 = device_loads(repl)
print(f"\nwith {R} redundant expert copies (EPLB-style)")
print(f" busiest GPU {d2.max():8.0f} tokens | mean {d2.mean():8.0f} | "
f"max/mean {d2.max()/d2.mean():.2f}")
print(f" idle time now {1 - d2.mean()/d2.max():.0%}")
# Communication: every token's activations travel to each chosen expert and back.
per_token = 2 * K * h * nbytes # dispatch + combine
print(f"\nall-to-all traffic per MoE layer")
print(f" per token : {per_token:>12,} bytes (2 x top-{K} x {h} x {nbytes})")
print(f" per batch of {T:,} : {per_token*T/2**30:>12.2f} GiB")
print(f" across 58 MoE layers: {per_token*T*58/2**30:>12.2f} GiB per forward pass")
# Node-limited routing bounds how many machines one token has to reach.
NODES = 8
node_of = torch.arange(E) % NODES
touched = torch.tensor([len(set(node_of[choice[t]].tolist())) for t in range(T)]).float()
print(f"\nnodes touched per token, unrestricted routing ({NODES} nodes)")
print(f" mean {touched.mean():.2f} max {int(touched.max())} "
f"(DeepSeek-V3 caps this at topk_group = 4)")256 experts on 32 GPUs, 8 experts each, 8,192 tokens, top-8 round-robin placement, no replication busiest GPU 3969 tokens | mean 2048 | max/mean 1.94 every GPU waits for the slowest -> 48% of GPU time is idle with 32 redundant expert copies (EPLB-style) busiest GPU 2932 tokens | mean 2048 | max/mean 1.43 idle time now 30% all-to-all traffic per MoE layer per token : 229,376 bytes (2 x top-8 x 7168 x 2) per batch of 8,192 : 1.75 GiB across 58 MoE layers: 101.50 GiB per forward pass nodes touched per token, unrestricted routing (8 nodes) mean 5.30 max 8 (DeepSeek-V3 caps this at topk_group = 4)
Reading the output
48% of GPU time is idle under plain round-robin placement. The busiest card holds 3,969 tokens' worth of work while the average is 2,048. Since the layer cannot finish until every card finishes, nearly half the cluster's capacity evaporates. This is the number that dominates real MoE serving economics.
Replicating 32 hot experts recovers a third of that. Max-to-mean drops from 1.94 to 1.43 and idle time from 48% to 30%, at the cost of 32 extra experts' worth of memory — about 1.4 GB at DeepSeek-V3's expert size. That is a very cheap trade, and it is why every serious MoE serving stack ships something like it.
101.5 GiB of all-to-all traffic per forward pass. For one batch of 8,192 tokens across 58 MoE layers. On a 400 Gb/s fabric that is roughly two seconds of pure network time if none of it is overlapped. Overlapping it with compute is not an optimisation, it is the difference between working and not.
5.30 nodes touched per token on average, out of 8. Almost every token needs almost every machine. DeepSeek-V3's topk_group: 4 caps this at 4 by construction, cutting cross-node traffic roughly in half before any kernel work.
The four kinds of parallelism, and which shortage each fixes
| Kind | What is split | Communication | Fixes |
|---|---|---|---|
| Data | the batch | all-reduce of gradients | throughput |
| Tensor | each weight matrix | all-reduce per layer | one layer too big |
| Pipeline | the layer stack | point-to-point activations | model too deep for one device |
| Expert | the experts | all-to-all per MoE layer | too many experts for one device |
| Context | the sequence | ring or all-gather of KV | sequence too long |
They compose. A production DeepSeek-V3 deployment typically runs expert parallelism across nodes and tensor parallelism within a node, with separate pools for prefill and decode.
What to actually use
- vLLM exposes expert-parallel deployment with EPLB, including window size and redundant expert count.
- SGLang integrates DeepSeek's EPLB. Their published large-scale deployment ran DeepSeek across 96 H100s with prefill/decode disaggregation and large-scale expert parallelism.
- DeepEP provides communication kernels tuned for MoE dispatch and combine, including low-latency paths for decoding.
- EPLB takes measured expert-load statistics and computes a placement, replicating hot experts across a configurable number of redundant slots. LPLB is a later research-stage variant that solves the assignment as a linear program per batch.
The pattern in all of them: measure real load, then place. A placement chosen from the architecture rather than from measurements will be wrong, because expert popularity depends on your traffic.
Common mistakes
Assuming uniform routing at serving time. Training balance is averaged over enormous batches. Your users send correlated requests, and per-batch imbalance is much worse. Measure it.
Enabling expert parallelism at small batch sizes. With few tokens per step, all-to-all latency dominates and there is nothing to overlap it with. Tensor parallelism is often better below a threshold you should measure.
Ignoring the combine step. Dispatch and combine are both all-to-all. Counting only dispatch halves your traffic estimate.
Running prefill and decode in one pool. Prefill is compute-bound with large token counts; decode is latency-bound with few. They want different batch sizes, different parallelism and different expert placements.
Setting redundant experts without measurement. Replicating the wrong experts costs memory and fixes nothing. Collect load counters first.
Try it yourself
Change torch.randn(E) * 0.8 to * 0.3 for a better-balanced router and see idle time fall without any replication. Then raise R to 64 and find where extra copies stop helping. That saturation point is the practical setting for the redundant-expert count.
What to learn next
- A working MoE layer in PyTorch — the layer these systems are distributing.
- Load balancing and expert collapse — the training-time half of the same problem.
- Model serving — the surrounding infrastructure this plugs into.
Researcher — Mathematics and papers.
The communication pattern
An MoE layer under expert parallelism is: all-to-all dispatch → local expert GEMMs → all-to-all combine. Per token per layer, bytes moved are
$$ 2 \, k \, d \, b $$
for top-$k$, hidden size $d$ and $b$ bytes per element, with the factor 2 counting dispatch and combine. For DeepSeek-V3 ($k=8$, $d=7168$, bf16) this is 229,376 bytes per token per layer, and 58 MoE layers give roughly 13.3 MB per token per forward pass.
Note that this is independent of the expert count $N$ and of expert width. Scaling to 512 experts costs nothing in all-to-all volume; it costs in the number of distinct destinations, which is a latency and scheduling problem rather than a bandwidth one.
The straggler bound
Layer time is set by the maximum device load, not the mean. With per-device loads $\ell_d$, utilisation is
$$ \eta = \frac{\bar{\ell}}{\max_d \ell_d} $$
The simulation above gives $\eta \approx 0.52$ for a moderately skewed router with one expert copy each. Published analyses of large MoE deployments report 20–40% of GPU cycles lost to hot-expert queuing, consistent with that range.
Expert popularity is heavy-tailed and, critically, stable over minutes. That stability is what makes measurement-driven placement work at all: a placement computed from the last few thousand batches remains good for the next few thousand.
Placement as an optimisation problem
DeepSeek's EPLB (Expert Parallelism Load Balancer) takes an expert-load vector and a device count, and returns an assignment that may replicate experts into redundant slots, minimising the maximum device load. With $R$ redundant slots, replicating expert $i$ across $r_i$ devices divides its load by $r_i$, and the problem is a makespan-minimisation over placements.
The greedy heuristic — replicate the top-$R$ hottest experts once each — is what the script implements and captures most of the available gain. LPLB formulates the token-to-replica assignment as a linear program solved per batch, dynamically reordering experts from workload statistics.
Two-level placement is standard in production: a hierarchical policy that first balances across nodes (minimising expensive cross-node traffic) and then within a node (where NVLink makes traffic cheap).
Node-limited routing
DeepSeek-V3 partitions experts into n_group: 8 groups aligned to physical nodes and restricts each token to topk_group: 4 of them. Cross-node destinations per token are bounded by 4 rather than by $\min(k, \text{nodes})$.
The simulation measures 5.30 nodes touched on average without the restriction. Capping at 4 is therefore a real reduction and, by the model's own reported ablations, close to free in quality.
This is the clearest case in current architectures of a modelling decision made for the interconnect. Expect more: at these expert counts, routing is a network scheduling problem.
Prefill and decode
The two phases have incompatible requirements.
Prefill has thousands of tokens per request, so per-expert token counts are large, GEMMs are efficient, and all-to-all is bandwidth-bound and overlappable. High expert-parallel degree pays.
Decode has one token per sequence per step. Per-expert counts are tiny, GEMMs are memory-bound, and all-to-all is latency-bound with nothing to hide behind. Low-latency dispatch kernels (DeepEP's decode path) exist precisely for this.
Disaggregated serving runs them in separate pools with separate parallel configurations, connected by KV cache transfer. LMSYS documented a DeepSeek deployment on 96 H100s built this way, with prefill/decode disaggregation and large-scale expert parallelism.
What to measure
- Per-device load, per layer, per step: max, mean and p99.
- All-to-all time as a fraction of layer time, dispatch and combine separately.
- Achieved overlap: layer time against the maximum of compute time and communication time.
- Expert popularity drift: how fast a placement computed an hour ago degrades.
- Cross-node against intra-node traffic, which is where routing restrictions show up.
Without the first and last of these, any expert-placement change is a guess.
What to learn next
- A working MoE layer in PyTorch — the layer these systems are distributing.
- Load balancing and expert collapse — the training-time half of the same problem.
- Model serving — the surrounding infrastructure this plugs into.