# How Expert Parallelism in MegaDLMs Optimizes MoE Model Training

> Discover how Expert Parallelism in MegaDLMs optimizes MoE model training. Reduce GPU memory and speed up training with efficient token routing across ranks. Learn more now.

- Repository: [Jinjie Ni/megadlms](https://github.com/jinjieni/megadlms)
- Tags: deep-dive
- Published: 2026-03-04

---

**Expert Parallelism (EP) in MegaDLMs partitions Mixture-of-Experts (MoE) layers across multiple GPUs via dedicated parallel groups, reducing per-GPU memory by distributing expert weights while using efficient All-to-All communication to route tokens between ranks.**

MegaDLMs extends NVIDIA's Megatron-LM architecture to support massive-scale Mixture-of-Experts models through **Expert Parallelism in MegaDLMs**. This specialized parallelism strategy creates a new dimension for distributing expert parameters across GPUs, enabling training with hundreds or thousands of experts without exhausting single-device memory limits.

## What Expert Parallelism Does in MegaDLMs

Expert Parallelism introduces a fourth parallel dimension alongside Tensor, Pipeline, and Data Parallelism. Instead of splitting individual tensors or layers, EP partitions the **set of experts** within an MoE layer across multiple GPUs.

| Parallelism Type | What It Splits | Runtime Representation |
|------------------|----------------|------------------------|
| **Tensor Parallel (TP)** | Individual tensor dimensions (e.g., weight matrix rows) | `_TENSOR_MODEL_PARALLEL_GROUP` |
| **Pipeline Parallel (PP)** | Whole transformer layers across stages | `_PIPELINE_MODEL_PARALLEL_GROUP` |
| **Data Parallel (DP)** | Full model replicas for gradient aggregation | `_DATA_PARALLEL_GROUP` |
| **Expert Parallel (EP)** | **Set of experts** inside an MoE layer | `_EXPERT_MODEL_PARALLEL_GROUP` |

The EP groups are initialized in `initialize_model_parallel` within [`megatron/core/parallel_state.py`](https://github.com/jinjieni/megadlms/blob/main/megatron/core/parallel_state.py). The `expert_model_parallel_size` parameter (defaulting to 1) determines how many ranks belong to each EP group. When this value exceeds 1, the **RankGenerator** constructs a parallel topology including the EP dimension via the `ep` token in `RankGenerator.__init__`.

## How Expert Parallelism Reduces Memory and Improves Compute

Expert Parallelism in MegaDLMs optimizes MoE training through three primary mechanisms: distributed expert storage, communication-aware token routing, and scalable device limiting.

### Memory Savings via Expert Sharding

Each GPU stores only the subset of experts assigned to its EP group. In [`megatron/core/transformer/moe/moe_layer.py`](https://github.com/jinjieni/megadlms/blob/main/megatron/core/transformer/moe/moe_layer.py), the MoE layer partitions expert weight tensors across EP ranks. This reduces per-GPU memory consumption to approximately `1 / expert_model_parallel_size` of the total MoE parameter count, enabling models with hundreds of experts to fit within standard GPU memory constraints.

### Efficient All-to-All Communication

The `MoEAlltoAllTokenDispatcher` in [`megatron/core/transformer/moe/token_dispatcher.py`](https://github.com/jinjieni/megadlms/blob/main/megatron/core/transformer/moe/token_dispatcher.py) orchestrates token routing across hybrid TP-EP groups:

1. **Gather Phase**: `token_permutation` calls `gather_from_sequence_parallel_region` on `self.tp_ep_group` to collect routing maps and probabilities across both Tensor and Expert Parallel groups.
2. **Expert Computation**: Local experts process their assigned token chunks using the partitioned weights stored only on their respective ranks.
3. **Scatter Phase**: `token_unpermutation` performs `reduce_scatter_to_sequence_parallel_region` to execute the reverse All-to-All, restoring the original token order while scaling outputs by expert probabilities.

This communication pattern minimizes cross-node traffic by confining token exchanges to the TP-EP hybrid groups defined in [`megatron/core/parallel_state.py`](https://github.com/jinjieni/megadlms/blob/main/megatron/core/parallel_state.py).

### Scalable Routing with Device Limiting

MegaDLMs supports `moe_router_topk_limited_devices` in `TransformerConfig` ([`megatron/core/transformer/transformer_config.py`](https://github.com/jinjieni/megadlms/blob/main/megatron/core/transformer/transformer_config.py)), which constrains how many EP ranks a token's routing decision considers. This reduces the volume of data gathered during the routing phase, improving latency for models with thousands of experts.

## Interaction with Other Parallelism Strategies

Expert Parallelism composes with existing Megatron-LM parallelism through carefully constructed hybrid groups.

### TP-EP Hybrid Groups

The framework builds combined Tensor-Expert parallel groups via `generator_wrapper('tp', is_expert=True)` and `generator_wrapper('tp-ep', is_expert=True)` in `initialize_model_parallel`. These hybrid communicators allow simultaneous tensor-parallel collectives and expert-parallel routing within the same process group, reducing communication overhead.

### Independence from Data Parallelism

EP groups operate orthogonally to Data Parallel groups. While DP gradients are reduced across full replicas in `_DATA_PARALLEL_GROUP`, expert-specific parameters are synchronized only within `_EXPERT_DATA_PARALLEL_GROUP`. This separation ensures that expert sharding does not interfere with standard data-parallel training loops.

## End-to-End Expert Parallelism Flow

Understanding the complete execution path clarifies how Expert Parallelism in MegaDLMs optimizes MoE models:

1. **Configuration**: `TransformerConfig` receives `num_moe_experts`, `expert_model_parallel_size`, and `moe_router_topk` parameters.
2. **Group Initialization**: `initialize_model_parallel` constructs EP groups and TP-EP hybrid groups in [`megatron/core/parallel_state.py`](https://github.com/jinjieni/megadlms/blob/main/megatron/core/parallel_state.py).
3. **Forward Pass**:
   - `MoELayer` invokes `MoEAlltoAllTokenDispatcher.token_permutation`, gathering hidden states and routing metadata across the TP-EP group.
   - Local experts (stored only on the ranks that own them) process their token chunk.
   - `token_unpermutation` performs the inverse All-to-All, scaling by expert probabilities and scattering back to the original sequence order.
4. **Backward Pass**: Gradients flow through the same communication paths, with expert parameter updates constrained to their respective EP groups.

## Code Examples

### Enabling Expert Parallelism via Command Line

Configure EP size alongside other parallel dimensions when launching training:

```bash

# Train a 64-expert MoE model on 8 GPUs with EP=2, TP=2, PP=2

python megadlms/main/tools/run_text_generation_server.py \
    --num-experts 64 \
    --expert-model-parallel-size 2 \
    --tensor-model-parallel-size 2 \
    --pipeline-model-parallel-size 2 \
    --moe-router-topk 8

```

The `--expert-model-parallel-size` flag creates the EP groups, while `--moe-router-topk` controls how many experts each token may access.

### Building an MoE Layer Programmatically

Instantiate a configured MoE layer that automatically uses Expert Parallelism:

```python
from megatron.core.transformer.transformer_config import TransformerConfig
from megatron.core.transformer.moe.moe_layer import MoELayer

config = TransformerConfig(
    hidden_size=4096,
    num_attention_heads=32,
    num_moe_experts=64,          # Total experts across all EP ranks

    expert_model_parallel_size=4, # 4-way EP → 16 experts per GPU

    moe_router_topk=8,            # Each token routed to top-8 experts

)

# Automatically partitions experts across EP groups

moe = MoELayer(config=config, layer_number=0)

```

The `MoELayer` internally creates a `MoEAlltoAllTokenDispatcher` that uses the EP communicator built in [`parallel_state.py`](https://github.com/jinjieni/megadlms/blob/main/parallel_state.py).

### Inspecting Expert Parallel Topology

Verify EP group configuration at runtime:

```python
from megatron.core.parallel_state import (
    get_expert_model_parallel_group,
    get_tensor_model_parallel_group,
)

ep_group = get_expert_model_parallel_group()
tp_group = get_tensor_model_parallel_group()

print(f"EP group world size: {ep_group.size()}")
print(f"TP group world size: {tp_group.size()}")

```

Typical output for an 8-GPU run with EP-size 2:

```

EP group world size: 2
TP group world size: 2

```

## Key Files in the MegaDLMs Repository

| File | Purpose |
|------|---------|
| [`megatron/core/parallel_state.py`](https://github.com/jinjieni/megadlms/blob/main/megatron/core/parallel_state.py) | Creates EP groups (`_EXPERT_MODEL_PARALLEL_GROUP`) and hybrid TP-EP groups via `initialize_model_parallel` and `RankGenerator`. |
| [`megatron/core/transformer/transformer_config.py`](https://github.com/jinjieni/megadlms/blob/main/megatron/core/transformer/transformer_config.py) | Defines `expert_model_parallel_size`, `moe_router_topk`, and `moe_router_topk_limited_devices` configuration parameters. |
| [`megatron/core/transformer/moe/token_dispatcher.py`](https://github.com/jinjieni/megadlms/blob/main/megatron/core/transformer/moe/token_dispatcher.py) | Implements `MoEAlltoAllTokenDispatcher` with `token_permutation` and `token_unpermutation` for A2A communication across TP-EP groups. |
| [`megatron/core/transformer/moe/moe_layer.py`](https://github.com/jinjieni/megadlms/blob/main/megatron/core/transformer/moe/moe_layer.py) | Defines `MoELayer` that orchestrates expert computation using EP-partitioned weights and the token dispatcher. |
| [`megatron/training/arguments.py`](https://github.com/jinjieni/megadlms/blob/main/megatron/training/arguments.py) | Exposes `--expert-model-parallel-size` CLI flag and documents parallelism incompatibilities. |
| [`megatron/core/transformer/moe/moe_utils.py`](https://github.com/jinjieni/megadlms/blob/main/megatron/core/transformer/moe/moe_utils.py) | Provides routing utilities including limited-device top-k selection for scalable EP. |
| [`tools/run_text_generation_server.py`](https://github.com/jinjieni/megadlms/blob/main/tools/run_text_generation_server.py) | Example entry point demonstrating EP argument propagation to the training loop. |

## Summary

- **Expert Parallelism in MegaDLMs** introduces a fourth parallel dimension that partitions MoE expert sets across GPUs, distinct from Tensor, Pipeline, and Data Parallelism.
- Memory consumption scales inversely with `expert_model_parallel_size`, allowing models with hundreds of experts to train on standard hardware by storing only `1/EP_size` of expert weights per GPU.
- The `MoEAlltoAllTokenDispatcher` orchestrates efficient token routing via `gather_from_sequence_parallel_region` and `reduce_scatter_to_sequence_parallel_region` across hybrid TP-EP groups, minimizing inter-GPU communication overhead.
- EP groups integrate seamlessly with existing Megatron-LM parallelism through `initialize_model_parallel` in [`megatron/core/parallel_state.py`](https://github.com/jinjieni/megadlms/blob/main/megatron/core/parallel_state.py), supporting combined TP-EP communicators and independent expert data-parallel reduction.

## Frequently Asked Questions

### What is Expert Parallelism and how does it differ from Tensor Parallelism in MegaDLMs?

**Expert Parallelism partitions the set of experts across GPUs, while Tensor Parallelism splits individual weight matrices within each expert.** In MegaDLMs, TP divides tensor dimensions like weight matrix rows across `_TENSOR_MODEL_PARALLEL_GROUP`, whereas EP distributes distinct expert instances across `_EXPERT_MODEL_PARALLEL_GROUP`. This means TP parallelizes the computation within a single expert, while EP parallelizes the number of experts available to the model.

### How does Expert Parallelism reduce memory usage for MoE models?

**EP reduces per-GPU memory by sharding expert parameters across the expert parallel group.** When `expert_model_parallel_size` is set to 4, each GPU in [`megatron/core/transformer/moe/moe_layer.py`](https://github.com/jinjieni/megadlms/blob/main/megatron/core/transformer/moe/moe_layer.py) stores only one-fourth of the total expert weights. This partitioning reduces memory consumption proportionally to the EP size, enabling training of models with 64 or 128 experts on hardware that could otherwise support only a fraction of that capacity.

### What communication operations does Expert Parallelism use to route tokens?

**EP relies on All-to-All communication through the `MoEAlltoAllTokenDispatcher`.** In [`megatron/core/transformer/moe/token_dispatcher.py`](https://github.com/jinjieni/megadlms/blob/main/megatron/core/transformer/moe/token_dispatcher.py), the `token_permutation` method performs `gather_from_sequence_parallel_region` across the TP-EP hybrid group to collect routing metadata. After local expert computation, `token_unpermutation` executes `reduce_scatter_to_sequence_parallel_region` to redistribute outputs. These operations ensure tokens reach their assigned experts while minimizing cross-node traffic within the defined parallel groups.

### Can Expert Parallelism be combined with Pipeline and Data Parallelism?

**Yes, EP composes orthogonally with Pipeline and Data Parallelism.** In [`megatron/core/parallel_state.py`](https://github.com/jinjieni/megadlms/blob/main/megatron/core/parallel_state.py), EP groups are constructed independently via `initialize_model_parallel`, allowing simultaneous use with PP stages and DP replicas. Expert-specific gradients are reduced within `_EXPERT_DATA_PARALLEL_GROUP`, while standard data-parallel reduction occurs across the full `_DATA_PARALLEL_GROUP`. This separation ensures that expert sharding does not interfere with existing training loops while supporting 4D parallelism configurations.