Serving a mixture-of-experts model on one node is mostly a configuration choice, covered in MoE serving architecture. Serving a model the size of DeepSeek-V3, with 256 routed experts in each of 58 MoE layers, is a deployment problem: the experts no longer fit comfortably on eight GPUs, so they are spread over two or more nodes, and every MoE layer of every forward pass now crosses the network twice.
This article is the runbook for that step. It sizes memory per GPU from first principles, launches a 16-GPU expert-parallel deployment across two nodes with vLLM, explains what the network must provide, and then deals with what makes multi-node EP operationally different from anything dense: the whole deployment is a single failure domain, it is only efficient at high concurrency, and its load changes as traffic changes which experts are hot.
The deployment shape
Wide expert parallelism combines two layouts. Attention runs data-parallel: each GPU is an independent rank with its own batch and its own KV cache, so attention needs no cross-GPU communication. The MoE layers run expert-parallel across all ranks: each GPU holds a slice of the experts, tokens are dispatched to the GPUs owning their chosen experts with an all-to-all, processed, and returned with a second all-to-all, called combine. With tensor-parallel size 1 the EP size equals the DP size, here 16.
The concepts behind dispatch and combine, and why decode all-to-alls are latency-bound, are in MoE all-to-all communication. What matters for deployment is the consequence: each MoE layer is a synchronisation point across every rank in the group, on every forward step.
Sizing memory per GPU
Work out the numbers before buying or reserving anything. DeepSeek-V3 has hidden size 7168 and expert intermediate size 2048, and each expert has three weight matrices (gate, up and down), so one expert in one layer has 3 x 7168 x 2048, about 44 million parameters, or about 44 MB in FP8. With 256 routed experts over 16 GPUs each GPU holds 16 experts per layer; across 58 MoE layers that is about 41 GB.
Every rank also holds a full copy of everything that is not a routed expert: attention, the shared expert, dense layers, embeddings and the output head. Treat that as roughly 17 GB in FP8 as a working estimate and measure it on your build. Add redundant experts for load balancing and the routed slice grows. The sizer below makes the arithmetic explicit and enforces the rule that physical experts, routed plus redundant, must divide evenly over the EP size:
def size_ep(n_routed, ep, redundant, moe_layers, hidden, inter, bytes_per_w,
dense_bytes, hbm_gb, util=0.90):
physical = n_routed + redundant
if physical % ep:
raise ValueError(f"{physical} physical experts do not divide over EP={ep}")
per_gpu = physical // ep
expert_bytes = 3 * hidden * inter * bytes_per_w # gate, up, down
experts_gb = per_gpu * expert_bytes * moe_layers / 1e9
weights_gb = experts_gb + dense_bytes / 1e9 # attention, shared expert, embeddings
kv_gb = hbm_gb * util - weights_gb
return per_gpu, round(experts_gb, 1), round(weights_gb, 1), round(kv_gb, 1)
# DeepSeek-V3 shape, FP8 weights, about 17 GB of non-routed weights per rank (estimate)
print(size_ep(256, 16, 16, 58, 7168, 2048, 1, 17e9, 80)) # (17, 43.4, 60.4, 11.6)
print(size_ep(256, 32, 32, 58, 7168, 2048, 1, 17e9, 80)) # (9, 23.0, 40.0, 32.0)On 80 GB GPUs at 90 percent utilisation, EP=16 with 16 redundant experts leaves about 12 GB per rank for KV cache; EP=32 leaves about 32 GB. Because attention is data-parallel, KV memory is per rank: a single long-context request must fit in one GPU's KV budget, not the pool's. That, more than raw throughput, is often what pushes teams to a wider EP group. Activations, CUDA graphs and communication buffers also take memory, so treat these numbers as an upper bound and confirm free KV blocks in the engine logs at start-up.
Network prerequisites
The all-to-all traffic is small per message at decode and large at prefill, and both kinds cross the node boundary. In-node traffic rides NVLink; cross-node traffic needs RDMA, typically InfiniBand or RoCE, with GPUDirect RDMA so the NICs read and write GPU memory without staging through the CPU. The DeepEP backends are built for exactly this and rely on NVSHMEM; check the DeepEP repository for the driver, NIC and library versions it currently requires rather than assuming.
Before launching the model, prove the fabric. Run an NCCL all-to-all benchmark across both nodes and confirm bandwidth close to the NIC line rate; confirm every GPU has an affine NIC on the same PCIe switch; and confirm the interface named in GLOO_SOCKET_IFNAME is the one the nodes can actually reach each other on. vLLM's guide calls this variable out because on InfiniBand clusters a wrong default interface makes initialisation hang rather than fail. A deployment that hangs at start-up is almost always a network naming or reachability problem, not a model problem.
Launching across two nodes
vLLM's multi-node EP launch uses one primary node, which runs the API servers and coordinates, and headless worker nodes that run only engine ranks. Every node passes the same total --data-parallel-size, its own --data-parallel-size-local, and the primary's address and RPC port; workers add --headless and a --data-parallel-start-rank equal to the number of ranks on the nodes before them.
# Node 0 (primary): API servers plus DP ranks 0-7
export GLOO_SOCKET_IFNAME=eth0 # the interface the nodes can reach each other on
vllm serve deepseek-ai/DeepSeek-V3 \
--tensor-parallel-size 1 \
--enable-expert-parallel \
--data-parallel-size 16 \
--data-parallel-size-local 8 \
--data-parallel-address 10.0.0.10 \
--data-parallel-rpc-port 13345 \
--api-server-count 8 \
--all2all-backend deepep_low_latency \
--enable-eplb \
--eplb-config '{"window_size": 1000, "step_interval": 3000, "num_redundant_experts": 16, "log_balancedness": true}'
# Node 1 (worker): no API server, ranks 8-15
export GLOO_SOCKET_IFNAME=eth0
vllm serve deepseek-ai/DeepSeek-V3 \
--headless \
--tensor-parallel-size 1 \
--enable-expert-parallel \
--data-parallel-size 16 \
--data-parallel-size-local 8 \
--data-parallel-start-rank 8 \
--data-parallel-address 10.0.0.10 \
--data-parallel-rpc-port 13345 \
--all2all-backend deepep_low_latency \
--enable-eplb \
--eplb-config '{"window_size": 1000, "step_interval": 3000, "num_redundant_experts": 16, "log_balancedness": true}'The backend choice follows the phase. deepep_low_latency is designed for decode, where per-step batches are small and latency dominates; deepep_high_throughput suits prefill, where messages are large; allgather_reducescatter is the portable default that works without DeepEP but scales worse. If you run one mixed deployment, choose for the phase that dominates your latency objective. If prefill and decode are disaggregated into separate deployments joined by KV transfer, each can use its own backend and even its own EP size. Keep the EPLB configuration identical on every node, and keep 16 + 256 divisible by 16, as the sizer checks.
Expert load balancing in operation
Routing is never uniform. Some experts receive far more tokens than others, and the hot set shifts with traffic: code prompts, one language or one customer can light up different experts. Because every MoE layer waits for the slowest rank, step time follows the most loaded GPU, not the average. EPLB records per-expert load over a sliding window of steps and periodically re-places experts, using the redundant slots to replicate hot ones on more than one GPU.
Turn on log_balancedness and watch the ratio of mean to maximum load per rank over time. If balancedness is poor right after a rebalance, the window may be too short and is chasing noise; if it degrades steadily between rebalances, the interval may be too long for how fast your traffic mix shifts. Rebalancing moves weights between GPUs, so it costs bandwidth and can cause latency blips; schedule-sensitive services should measure p99 latency around rebalance events before tightening the interval. MoE serving architecture explains the max-not-mean problem in more depth.
One EP group is one failure domain
In a dense data-parallel fleet, losing a GPU loses one replica's worth of capacity. In wide EP, every rank holds experts that every other rank needs, so a rank that crashes, hangs or slows down stalls every MoE layer on every rank. The practical consequences:
- Health is group-level. The load balancer should route away from the whole deployment when any rank is unhealthy, and the readiness check should exercise a real forward pass, not just the HTTP port.
- Restart is group-level. Plan to restart the whole group, and budget the start-up time: loading hundreds of gigabytes of weights, warming kernels and capturing graphs takes minutes.
- Capacity comes in groups. Run at least two independent EP groups behind the load balancer for anything with an availability target, so one can drain or restart while the other serves.
- Stragglers are failures. A GPU with thermal throttling or a NIC with link errors does not crash, it slows every step. Alert on per-rank step time spread, not just on errors.
On Kubernetes, model each EP group as one unit that is scheduled, started and replaced together, for example with a leader-and-workers pattern, and place its pods within one high-bandwidth fabric domain.
Smoke test and rollout
Never put a new EP group behind production traffic on the strength of a successful start. Run a smoke test that sends concurrent, deterministic requests, checks every response is non-empty, and records latency percentiles. Concurrency matters: EP bugs often appear only when many ranks have work at once.
import asyncio, aiohttp, time
PROMPTS = ["Explain TCP slow start.", "Write a haiku about GPUs.", "Sum 1..100 in Python."] * 64
async def one(session, url, prompt):
t0 = time.perf_counter()
async with session.post(url, json={"model": "deepseek-ai/DeepSeek-V3", "prompt": prompt,
"max_tokens": 128, "temperature": 0}) as r:
if r.status != 200:
return r.status, time.perf_counter() - t0, ""
body = await r.json()
return r.status, time.perf_counter() - t0, body["choices"][0]["text"]
async def smoke(url="http://10.0.0.10:8000/v1/completions"):
async with aiohttp.ClientSession() as s:
res = await asyncio.gather(*(one(s, url, p) for p in PROMPTS))
ok = [r for r in res if r[0] == 200 and r[2].strip()]
lat = sorted(r[1] for r in ok)
assert len(ok) == len(PROMPTS), f"{len(PROMPTS) - len(ok)} failed or empty"
print(f"p50={lat[len(lat)//2]:.2f}s p99={lat[int(len(lat)*0.99)]:.2f}s")
asyncio.run(smoke())Roll out by group. Bring up the new version as a fresh EP group, run the smoke test and a short quality check against a fixed evaluation set at temperature zero, then shift a small share of traffic, compare p50 and p99 time-to-first-token and inter-token latency against the old group, and only then drain the old group. Rolling individual ranks to a new version inside a running group is not a supported pattern; ranks must agree on the model, the expert layout and the backend.
Worked example: choosing EP=16 or EP=32
A team serving DeepSeek-V3 with a 64k-token context target compares two 80 GB layouts. At EP=16 the sizer leaves about 12 GB of KV per rank. MLA keeps the KV cache small per token, but 64k tokens per request, times enough concurrent requests per rank to make the MoE layers efficient, does not fit, so they would be forced to cap concurrency, which is exactly what makes MoE inefficient. At EP=32 each rank holds 9 experts per layer and has about 32 GB of KV, enough for the target concurrency.
The cost is a second pair of nodes per group and more cross-node traffic: with 32 ranks over four nodes, roughly three quarters of each token's eight expert destinations are on other nodes, against about half at EP=16. They choose EP=32 for decode, accept the higher per-group cost, and keep two groups for availability. The arithmetic behind the communication volume is in EP math.
Failure modes
| Symptom | Likely cause | Fix |
|---|---|---|
| Start-up hangs at initialisation | Wrong interface for the control plane, or ports blocked between nodes | Set GLOO_SOCKET_IFNAME; test reachability on the RPC port |
| Starts, then out-of-memory on first long request | KV budget per rank smaller than one request | Cap max model length or widen EP |
| Error on start about expert counts | Routed plus redundant experts not divisible by EP size | Adjust num_redundant_experts |
| Throughput far below expectation | Low concurrency, or cross-node traffic not using RDMA | Load test at target concurrency; verify GPUDirect RDMA in the NCCL benchmark |
| p99 spikes every few minutes | EPLB rebalancing or a straggling rank | Correlate with rebalance logs and per-rank step time |
| Whole service down after one GPU fault | Single EP group behind the load balancer | Run two or more groups and health-check at group level |
What to do next
- Compute per-GPU expert, dense and KV memory for your model with the sizer, and check physical experts divide evenly over EP.
- Benchmark the cross-node all-to-all with NCCL and confirm RDMA is in use before launching the model.
- Launch one primary and headless workers with identical data-parallel, backend and EPLB settings.
- Pick the all-to-all backend for the phase that dominates your latency target, or disaggregate prefill and decode.
- Enable balancedness logging and correlate it with latency around rebalances.
- Run at least two EP groups, with readiness checks that execute a real forward pass.
- Roll out by group with a concurrent smoke test and a fixed-evaluation quality check.