How Expert Parallelism in MegaDLMs Optimizes MoE Model Training
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. 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, 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 orchestrates token routing across hybrid TP-EP groups:
- Gather Phase:
token_permutationcallsgather_from_sequence_parallel_regiononself.tp_ep_groupto collect routing maps and probabilities across both Tensor and Expert Parallel groups. - Expert Computation: Local experts process their assigned token chunks using the partitioned weights stored only on their respective ranks.
- Scatter Phase:
token_unpermutationperformsreduce_scatter_to_sequence_parallel_regionto 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.
Scalable Routing with Device Limiting
MegaDLMs supports moe_router_topk_limited_devices in TransformerConfig (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:
- Configuration:
TransformerConfigreceivesnum_moe_experts,expert_model_parallel_size, andmoe_router_topkparameters. - Group Initialization:
initialize_model_parallelconstructs EP groups and TP-EP hybrid groups inmegatron/core/parallel_state.py. - Forward Pass:
MoELayerinvokesMoEAlltoAllTokenDispatcher.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_unpermutationperforms the inverse All-to-All, scaling by expert probabilities and scattering back to the original sequence order.
- 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:
# 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:
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.
Inspecting Expert Parallel Topology
Verify EP group configuration at runtime:
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 |
Creates EP groups (_EXPERT_MODEL_PARALLEL_GROUP) and hybrid TP-EP groups via initialize_model_parallel and RankGenerator. |
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 |
Implements MoEAlltoAllTokenDispatcher with token_permutation and token_unpermutation for A2A communication across TP-EP groups. |
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 |
Exposes --expert-model-parallel-size CLI flag and documents parallelism incompatibilities. |
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 |
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 only1/EP_sizeof expert weights per GPU. - The
MoEAlltoAllTokenDispatcherorchestrates efficient token routing viagather_from_sequence_parallel_regionandreduce_scatter_to_sequence_parallel_regionacross hybrid TP-EP groups, minimizing inter-GPU communication overhead. - EP groups integrate seamlessly with existing Megatron-LM parallelism through
initialize_model_parallelinmegatron/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 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, 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, 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.
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 →