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_eplb and eplb_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:

  1. State Tracking: EplbState maintains per-expert call statistics across the EP group.
  2. Remapping Decision: The algorithm identifies hot experts and redistributes them using either "linear" or "round_robin" placement strategies.
  3. 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:

  1. Initialization: ParallelConfig.enable_expert_parallel triggers creation of the _EP group in parallel_state.py.
  2. Model Sharding: MoEMixin.recursive_replace substitutes experts ModuleLists with TransformersFusedMoE instances that store ep_size and num_local_physical_experts.
  3. Routing: The router computes local topk_ids, then gathers them across the EP group (or DP group if sequence parallelism is enabled).
  4. Token Dispatch: The fused kernel performs all-to-all communication via the selected all2all_backend to route hidden states to expert owners.
  5. Local Computation: Each GPU runs its assigned experts on received token slices.
  6. Result Return: A second all-to-all aggregates expert outputs back to the token originators.
  7. 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

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:

Share the following with your agent to get started:
curl -s "https://instagit.com/install.md"

Works with
Claude Codex Cursor VS Code OpenClaw Any MCP Client

Maintain an open-source project? Get it listed too →