How Expert Parallelism Works for Mixture-of-Experts Models in vLLM
Expert parallelism (EP) in vLLM creates a dedicated data-parallel process group that shards expert weight matrices across GPUs, replacing tensor parallelism for MoE layers to enable scalable distributed inference.
Expert parallelism is a critical optimization for deploying large Mixture-of-Experts (MoE) language models across multiple GPUs in the vLLM inference engine. Unlike traditional tensor parallelism, which shards every layer uniformly, EP specifically distributes the heavy feed-forward expert networks to reduce memory pressure and minimize all-reduce communication overhead. This article explains the complete implementation of EP in vLLM, from configuration flags in ParallelConfig to the all-to-all communication kernels in the MoE forward pass.
Configuring Expert Parallelism
Users enable EP through ParallelConfig in vllm/config/parallel.py. The system activates only when loading a model with num_experts > 0.
enable_expert_parallel: Boolean flag that switches MoE layers from tensor parallelism to EP sharding.expert_placement_strategy: Determines expert distribution across ranks, accepting"linear"or"round_robin".enable_eplbandeplb_config: Optional Expert Parallel Load Balancing that periodically reshuffles experts to improve GPU utilization.all2all_backend: Selects the communication backend (e.g., NCCL) for the all-to-all token exchange between EP ranks.
These parameters determine how vllm.distributed.parallel_state initializes the expert-parallel world size and communication patterns.
EP Process Group Initialization
When a MoE model is launched, vLLM creates a dedicated expert-parallel process group (_EP) that operates independently of the standard data-parallel and tensor-parallel groups. The entry point is the get_ep_group() function in vllm/distributed/parallel_state.py:
# vllm/distributed/parallel_state.py
def get_ep_group() -> GroupCoordinator:
assert _EP is not None, (
"expert parallel group is not initialized. "
"EP group is only created for MoE models with num_experts > 0. "
"This function should only be called for MoE models."
)
return _EP
The group size (ep_size) equals the world size of _EP and determines how many experts each rank owns. This group is strictly isolated from tensor-parallel groups when EP is enabled.
MoE Layer Construction and Local Expert Sharding
The MoEMixin class in vllm/model_executor/models/transformers/moe.py handles the conversion of standard MoE layers into expert-parallel variants. During model construction, it calculates the local expert count:
# vllm/model_executor/models/transformers/moe.py
self.num_physical_experts = num_experts + num_redundant_experts
self.num_local_physical_experts = self.num_physical_experts // ep_size
Here, ep_size = get_ep_group().world_size replaces the tensor-parallel size when enable_expert_parallel is True. Each rank loads only its assigned num_local_physical_experts weights into memory, reducing per-GPU footprint proportionally to the EP world size.
Token Routing and All-to-All Communication
During the forward pass, the router selects top-k experts for each token. In EP mode, the routing indices (topk_ids) must be visible across the entire EP group before the token exchange. The custom_routing_function in moe.py gathers these IDs using all-gather operations:
# vllm/model_executor/models/transformers/moe.py
if topk_ids.size(0) != hidden_states.size(0):
dp_metadata = get_forward_context().dp_metadata
sizes = dp_metadata.get_chunk_sizes_across_dp_rank()
is_sp = self.is_sequence_parallel
dist_group = get_ep_group() if is_sp else get_dp_group()
(topk_ids,) = dist_group.all_gatherv([topk_ids], 0, sizes)
After gathering, the fused MoE kernel in vllm/model_executor/layers/fused_moe/fused_moe.py executes the all-to-all communication (using the configured all2all_backend) to move hidden states to the target EP ranks. Each rank processes tokens only for its local experts, then a reverse all-to-all returns results to the original token positions.
Expert Parallel Load Balancing (EPLB)
When enable_eplb is active, vLLM monitors expert utilization through a sliding window defined by EPLBConfig.window_size. After every EPLBConfig.step_interval forward steps, the EPLB routine in vllm/distributed/eplb/eplb_state.py remaps logical experts to physical GPU locations:
- State Tracking:
EplbStatemaintains per-expert call statistics across the EP group. - Remapping Decision: The algorithm identifies hot experts and redistributes them using either
"linear"or"round_robin"placement strategies. - Async Migration: When configured, expert weights migrate between ranks asynchronously to minimize pipeline bubbles.
This dynamic rebalancing prevents GPU starvation when token routing patterns are skewed.
End-to-End Execution Flow
The complete EP lifecycle in vLLM follows this pipeline:
- Initialization:
ParallelConfig.enable_expert_paralleltriggers creation of the_EPgroup inparallel_state.py. - Model Sharding:
MoEMixin.recursive_replacesubstitutesexpertsModuleLists withTransformersFusedMoEinstances that storeep_sizeandnum_local_physical_experts. - Routing: The router computes local
topk_ids, then gathers them across the EP group (or DP group if sequence parallelism is enabled). - Token Dispatch: The fused kernel performs all-to-all communication via the selected
all2all_backendto route hidden states to expert owners. - Local Computation: Each GPU runs its assigned experts on received token slices.
- Result Return: A second all-to-all aggregates expert outputs back to the token originators.
- Load Balancing: Optional EPLB periodically reshuffles expert placement based on utilization statistics.
Code Examples
Enabling Expert Parallelism via CLI
Launch a 4-GPU inference server with EP enabled (2 data-parallel ranks, experts sharded across them):
python -m vllm.entrypoints.openai.api_server \
--model mixtral-8x7b \
--tensor-parallel-size 1 \
--data-parallel-size 2 \
--enable-expert-parallel \
--expert-placement-strategy round_robin
Configuring EP in Python
from vllm import LLM, SamplingParams
llm = LLM(
model="mixtral-8x7b",
enable_expert_parallel=True,
expert_placement_strategy="linear",
enable_eplb=True,
eplb_config=dict(num_redundant_experts=2)
)
sampling_params = SamplingParams(temperature=0.7, top_p=0.9)
outputs = llm.generate(
prompts=["Explain expert parallelism in vLLM."],
sampling_params=sampling_params
)
print(outputs[0].text)
Accessing EP Metadata at Runtime
Inspect the expert-parallel world size and rank during execution:
from vllm.distributed import get_ep_group
ep_group = get_ep_group()
print(f"EP size: {ep_group.world_size}, rank: {ep_group.rank}")
Summary
- Expert parallelism creates a distinct process group (
_EP) invllm/distributed/parallel_state.pythat shards expert weights separately from tensor parallelism. - Configuration occurs through
ParallelConfigflags invllm/config/parallel.py, including optional load balancing viaEPLBConfig. - The
MoEMixinclass invllm/model_executor/models/transformers/moe.pycalculatesnum_local_physical_expertsby dividing total experts byep_size. - Token routing requires gathering
topk_idsacross the EP group before executing the all-to-all dispatch in the fused MoE kernel. - EPLB dynamically redistributes experts based on utilization statistics stored in
vllm/distributed/eplb/eplb_state.py.
Frequently Asked Questions
What is the difference between expert parallelism and tensor parallelism in vLLM?
Tensor parallelism (TP) shards every layer's weights (attention and feed-forward) across GPUs, requiring all-reduce operations after each layer. Expert parallelism (EP) only shards the MoE feed-forward networks, keeping attention layers fully replicated or using standard TP. EP reduces communication volume for MoE models because experts are sparse; only activated tokens move via all-to-all, while TP all-reduces the full activation tensor for every layer.
When should I enable expert parallelism instead of standard tensor parallelism?
Enable EP when serving large MoE models (e.g., Mixtral 8x7B, DeepSeek-V2) where expert parameters dominate memory usage. EP is preferable when the number of experts is large relative to the tensor-parallel size, as it partitions the heavy expert matrices instead of replicating them on every TP rank. If your deployment uses data parallelism with multiple replicas, EP combined with a smaller TP size often yields better throughput than pure TP scaling.
How does the all-to-all backend affect expert parallelism performance?
The all2all_backend parameter in ParallelConfig controls the communication library used for token dispatch. The default NCCL backend typically provides optimal latency for GPU-to-GPU transfers within a node. For cross-node EP deployments, selecting the appropriate backend minimizes the latency of moving hidden states to remote experts during the forward pass, directly impacting end-to-end inference latency.
Does enabling EPLB (Expert Parallel Load Balancing) impact inference latency?
EPLB introduces periodic remapping of experts based on sliding window statistics, which requires additional communication to migrate expert weights between ranks. However, vLLM implements EPLB asynchronously when possible, overlapping the reshuffling with computation to hide latency. The benefit of improved load balancing—preventing GPU starvation on hot experts—typically outweighs the brief overhead of periodic remapping in high-throughput serving scenarios.
Have a question about this repo?
These articles cover the highlights, but your codebase questions are specific. Give your agent direct access to the source. Share this with your agent to get started:
curl -s "https://instagit.com/install.md" Maintain an open-source project? Get it listed too →