# How Expert Parallelism Works for Mixture-of-Experts Models in vLLM

> Discover how expert parallelism in vLLM shards expert weights across GPUs for scalable distributed inference of Mixture-of-Experts models. Optimize your MoE inference.

- Repository: [vLLM/vllm](https://github.com/vllm-project/vllm)
- Tags: deep-dive
- Published: 2026-03-03

---

**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`](https://github.com/vllm-project/vllm/blob/main/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`](https://github.com/vllm-project/vllm/blob/main/vllm/distributed/parallel_state.py):

```python

# 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`](https://github.com/vllm-project/vllm/blob/main/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:

```python

# 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`](https://github.com/vllm-project/vllm/blob/main/moe.py) gathers these IDs using all-gather operations:

```python

# 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`](https://github.com/vllm-project/vllm/blob/main/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`](https://github.com/vllm-project/vllm/blob/main/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`](https://github.com/vllm-project/vllm/blob/main/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):

```bash
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

```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:

```python
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`) in [`vllm/distributed/parallel_state.py`](https://github.com/vllm-project/vllm/blob/main/vllm/distributed/parallel_state.py) that shards expert weights separately from tensor parallelism.
- Configuration occurs through `ParallelConfig` flags in [`vllm/config/parallel.py`](https://github.com/vllm-project/vllm/blob/main/vllm/config/parallel.py), including optional load balancing via `EPLBConfig`.
- The `MoEMixin` class in [`vllm/model_executor/models/transformers/moe.py`](https://github.com/vllm-project/vllm/blob/main/vllm/model_executor/models/transformers/moe.py) calculates `num_local_physical_experts` by dividing total experts by `ep_size`.
- Token routing requires gathering `topk_ids` across 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`](https://github.com/vllm-project/vllm/blob/main/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.