Mixture-of-experts serving
A mixture-of-experts model replaces a single dense feed-forward block with many smaller ones and a router that picks which few run for each token. DeepSeek-V3 is a concrete example of the shape this takes: each MoE layer has 1 shared expert, always active, plus 256 routed experts, of which 8 are activated for any given token through expert routing, a sigmoid affinity score per expert with a per-expert bias term adjusted over training to keep load balanced without an auxiliary loss that would otherwise hurt model quality. The result is a model with 671B total parameters where only 37B are activated per token, which is the entire economic argument for MoE: most of the parameters sit idle for any single forward pass.
That sparsity has to be realized physically, and that is where expert parallelism and its communication pattern take over. DeepSeek-V3’s prefilling deployment unit spans 4 nodes and 32 GPUs, combining 4-way tensor parallelism for attention with 32-way expert parallelism for the MoE blocks, so that “each expert processes a sufficiently large batch size.” Because a token’s chosen experts can live on any GPU in that group, every MoE layer needs an all-to-all exchange: each GPU sends token activations out to the GPUs hosting its routed experts (dispatch) and receives the computed results back (combine). DeepEP, DeepSeek’s open-source dispatch and combine library, is that exchange packaged: it provides “high-throughput and low-latency all-to-all GPU kernels (MoE dispatch and combine) with low-precision support including FP8,” compiled at runtime by a just-in-time module so installation needs no CUDA compilation. Its requirements name both fabrics separately, NVLink for intranode communication and an RDMA network for internode communication, because a single dispatch uses both.
Because the router decides per token which experts get used, expert load is never guaranteed even, and an overloaded expert stalls the GPU that hosts it while the rest of the all-to-all waits. DeepSeek-V3 addresses this by deploying redundant copies of the experts that online traffic statistics show are hottest, rebalancing that set roughly every 10 minutes; its decoding deployment goes further, spanning 40 nodes and 320 GPUs with one expert per GPU and 64 GPUs dedicated to hosting redundant and shared experts, using direct point-to-point transfers over InfiniBand to keep dispatch and combine latency low. The load-balancing problem and the communication problem are really the same problem seen from two sides: a router that is free to send every token anywhere needs both a placement strategy that keeps GPUs evenly loaded and a communication kernel fast enough that the resulting all-to-all does not dominate the layer’s runtime.
DeepSeek open-sourced the placement half of that answer as EPLB, its expert parallelism load balancer, which computes “a balanced expert replication and placement plan based on the estimated expert loads”; predicting those loads is left to the deployer, with a moving average of historical statistics named as the common method. It ships two policies. Hierarchical load balancing, meant for the prefilling stage with a smaller expert-parallel size, packs whole expert groups onto nodes evenly and then replicates within each node, exploiting DeepSeek-V3’s group-limited routing to place “the experts of the same group to the same node to reduce inter-node data traffic.” Global load balancing, for the decoding stage with a larger expert-parallel size, replicates experts across the whole fleet regardless of groups. Either way, replication and placement fall out of measured load rather than being fixed at deployment: the plan is an output of traffic.
MegaScale-Infer pushes the same disaggregation logic inside the MoE layer itself. Its starting observation is that sparse activation “shifts feed-forward networks (FFNs) from being compute-intensive to memory-intensive during inference, leading to substantially lower GPU utilization,” because each expert sees only its routed slice of the batch. The system therefore disaggregates the attention and FFN modules within each layer onto separate GPUs, a strategy it names disaggregated expert parallelism, and then gives each module the parallelism that suits it: “attention modules are replicated using data parallelism, while FFN modules are scaled with expert parallelism,” so that consolidating the requests of many attention replicas is what pushes each expert’s batch back up. On top of that it runs ping-pong pipeline parallelism: a request batch is partitioned into micro-batches that shuttle between the attention side and the expert side, so one side computes while the other’s traffic is in flight. Evaluated on MoE models from 132 to 317 billion parameters, MegaScale-Infer reports up to 1.90× higher per-GPU throughput than state-of-the-art systems, and on a heterogeneous cluster, where attention sits on GPUs bought for memory and the FFNs on GPUs bought for compute, 1.7× higher throughput per dollar.
That system also carries the sharpest counterexample to this chapter’s own framing of NCCL as the way GPUs exchange data. Once attention and FFN each carry their own parallelism configuration, the token routing between them stops being an all-to-all among equals and becomes M2N communication, “where M and N represent the number of senders and receivers, respectively.” The authors report “performance shortcomings of popular communication libraries” on that specific pattern and write their own, one that eliminates “unnecessary GPU-to-CPU data copies, group initialization overhead, and GPU synchronization.” The measurement they publish is not marginal: “Compared to NCCL, a widely-used communication library, MegaScale-Infer’s M2N communication achieves 4.2× higher throughput and 68.2% lower latency.” NCCL is the right tool when the participants form a symmetric group all running the same collective, which is exactly what tensor and pipeline parallelism ask of it. It stops being the obvious tool when the senders and the receivers are different in number, on different machines, and under different parallel strategies. Generality is priced, and MoE serving is where the chapter’s default library starts paying for it.
The all-to-all as a kernel#
Underneath the deployment units and the placement plans, dispatch and combine are kernels, and DeepSeek’s own account of how they are written is the closest this chapter comes to the rest of the book. The V3 report describes cross-node all-to-all kernels co-designed with the MoE gating algorithm and with the network topology of the cluster they ran on, where “NVLink offers a bandwidth of 160 GB/s, roughly 3.2 times that of IB (50 GB/s).” Everything else follows from that ratio. Each token is limited to at most 4 nodes, which caps its InfiniBand traffic. Once its routing decision is made it crosses IB “to the GPUs with the same in-node index on its target nodes,” and only from there is forwarded over NVLink to whichever GPU actually holds its expert. Addressing that hop to a fixed in-node index rather than to the destination GPU is what decouples the two transfers: IB and NVLink communication “are fully overlapped” instead of running one after the other. The payoff is stated as headroom rather than as a speedup. Because a token can then average 3.2 experts per node without extra NVLink cost, the model “can scale up this number to a maximum of 13 experts (4 nodes × 3.2 experts/node) while preserving the same communication cost,” even though V3 routes to 8. The router’s top-k is bounded by the network, and the kernel is what sets the bound.
The rest of that section reads like this book’s kernel chapters applied to communication. Only 20 SMs are “sufficient to fully utilize the bandwidths of IB and NVLink,” and those 20 are partitioned into 10 communication channels using warp specialization, the same division of a block’s warps into roles that kernel optimization uses to keep tensor cores fed from a circular buffer. Here the roles are transport stages. During dispatch, IB sending, IB-to-NVLink forwarding, and NVLink receiving each get their own warps, and the number allocated to each task is “dynamically adjusted according to the actual workload across all SMs”; during combine, the mirror image, NVLink sending, NVLink-to-IB forwarding and accumulation, and IB receiving and accumulation. Because both kernels overlap the computation stream, their cache footprint is part of the design rather than an afterthought: the report uses “customized PTX (Parallel Thread Execution) instructions” and auto-tunes the communication chunk size, which “significantly reduces the use of the L2 cache and the interference to other SMs.” The report presents all of this in its training infrastructure section, where the kernels hide inside a pipeline schedule, but dispatch and combine are the same exchange a serving deployment runs at every MoE layer. A communication library that costs 20 SMs and stays out of the L2 is not a networking result. It is a kernel result.
The open-source library has since moved past the version that shipped with V3, in a direction that makes the SM budget the headline. DeepEP’s V2 release is a complete refactoring of expert parallelism, “achieving extreme performance with several times fewer SM resources compared to V1, while supporting significantly larger scale-up and scale-out domains,” and it “has also switched from the NVSHMEM backend to the more lightweight NCCL Gin backend.” The separate high-throughput and low-latency APIs are unified behind a single ElasticBuffer interface with a new GEMM layout; the domains reach EP2048; SM and queue-pair counts are now calculated analytically, with “no more auto-tuning needed”; and for V3-like legacy training, “SM usage reduced from 24 to 4 - 6 while maintaining equivalent or better performance.” Following V3’s configuration, 8K tokens per batch, 7168 hidden dimensions, top 8 experts, FP8 dispatching and BF16 combining, the README reports an SM100 machine on CX7 NICs at EP 8 x 2 reaching 90 GB/s dispatch and 91 GB/s combine bottleneck bandwidth on 12 SMs, and claims up to 1.3× V1’s peak performance while saving up to 4× the SM count. The 20-SM design above is what was published with the model; the library it became asks for fewer. When a chapter cites a communication library, the version is part of the citation.
Turning expert parallelism on#
What DeepSeek describes as a deployment unit, a serving engine exposes as one flag and an identity. vLLM’s expert parallelism “allows experts in Mixture-of-Experts (MoE) models to be deployed on separate GPUs, increasing locality, efficiency, and throughput overall,” and it is not sized directly: the expert parallel size is computed as the tensor parallel size multiplied by the data parallel size. Turning it on means picking those two numbers and setting one boolean.
# Single node EP deployment
vllm serve deepseek-ai/DeepSeek-V3-0324 \
--tensor-parallel-size 1 \
--data-parallel-size 8 \
--enable-expert-parallelOne times eight is eight, so that command builds an expert-parallel group spanning a whole node, and the documentation notes it fits DeepSeek-V3-0324 on an eight-GPU H200 or H20. The two kinds of layer are then treated differently, which is the whole reason to set the flag: the expert layers are sharded across all eight ranks, while the attention weights, at tensor parallel size 1, are replicated on every one of them. Raise the tensor parallel size above one and attention is sharded across those ranks inside each data parallel group instead. Drop the flag entirely and the MoE layers fall back to tensor parallelism over a group of size eight, the same treatment a dense model gets. DeepSeek-V3’s prefilling unit, tensor parallelism for attention and expert parallelism for the MoE blocks, is those two rules applied to one deployment.
Past one node the shape holds and the launch splits. Each machine runs its own command carrying the total data parallel size and its own local share; the first node serves the API and every other node starts headless, so all client requests are handled by the primary and the others are workers only. What stitches the ranks together is a start-rank offset on each secondary node equal to the cumulative local size of the nodes before it. The model is still one deployment. The command is now one per machine.
The flag is cheap; what has to exist underneath it is not. vLLM’s prerequisites for expert parallelism are DeepEP built against the host environment, the DeepGEMM library, and, for disaggregated serving, gdrcopy, with its newer DeepEP backend requiring NCCL 2.30.4 or later because the version PyTorch ships is older. Which all-to-all kernel then runs is a separate choice: --all2all-backend defaults to an implementation built from allgather and reduce-scatter primitives that works with any expert- and data-parallel configuration, and the two DeepEP modes are split by phase, a high-throughput one for multi-node prefill using grouped GEMM with a continuous layout and a low-latency one for multi-node decode with CUDA graph support and a masked layout. The prefill and decode distinction the next section draws between machines shows up here, one level down, as which kernel carries the dispatch.
EPLB arrives in the same place, as a pair of flags. --enable-eplb makes vLLM collect load statistics on every forward pass and rebalance the expert mapping periodically, and --eplb-config carries the schedule: a window of engine steps to track for the decision, 1000 by default; an interval between rebalances, 3000 steps by default; a count of redundant experts per rank beyond the equal split; and a switch that logs balancedness, defined as the average tokens per expert divided by the maximum. DeepSeek’s ten-minute rebalancing cadence is the same knob expressed in wall-clock rather than engine steps, and its redundant hot experts are that count. The price is stated plainly: the extra copies have to fit in GPU memory, roughly 2.4 GB for DeepSeek-V3 at one redundant expert per rank, so vLLM warns that EPLB “may not be a good fit for memory constrained environments or when KV cache space is at a premium.” Expert balance is bought with cache.