MoonshotAI/MoonEP
GitHub: MoonshotAI/MoonEP
MoonEP 是一个面向大规模 MoE 模型的专家并行通信库,通过动态冗余专家机制实现各 rank 间 token 负载的完美均衡,解决分布式训练中路由不平衡导致的性能退化和显存碎片问题。
Stars: 80 | Forks: 7
# MoonEP
MoonEP 是一个专家并行通信库,它通过动态冗余专家来保持各个 rank 间 token 负载的完美均衡。
**符号说明**:`S` = 每个 rank 的输入 token 数,`K` = 每个 token 路由的 top-k。
1. **完美均衡**:无论路由倾斜程度如何,每个 rank 都会精确接收 `S × K` 个 token。少量冗余专家会根据当前路由输出进行在线规划,并在专家计算前预取;在反向传播中,它们的梯度会被归约回其原始 rank。
2. **在线规划**:近乎最优的 GPU 规划 kernel,开销微乎其微。
3. **零拷贝与静态 shape**:融合的 permute/unpermute —— token 被直接发送到远程 rank 上按专家分组的位置,并将 buffer 视图返回给计算过程。仅需要一个固定的 `S × K` buffer,并且静态已知的 shape 消除了逐层 MoE 的主机同步。
## 性能
这两项基准测试均在 EP=8 的 H20 上运行,并不断调整路由不平衡度:
$$\text{maxvio} = \max_e \left( \frac{T_e}{\bar{T}} \right) - 1$$
其中 $T_e$ 是路由到专家 $e$ 的 token 数量,而 $\bar{T}$ 是在完美均衡下(maxvio = 0 表示完美均衡)每个专家预期的 token 数量。
**通信对比 DeepEP v2** ([benchmarks/bench_vs_deepep.py](benchmarks/bench_vs_deepep.py)):
- **零拷贝使原始通信更快**:token 被直接写入到它们在远程 rank 上最终的按专家分组的位置 —— 无需 permute 输入,也无需 permute 输出 —— 通信 buffer 的视图被直接交给计算过程,从而消除了在尾声阶段占主导地位的 通信 buffer → 用户 buffer 拷贝。MoonEP 的通信时间在每种不平衡级别下都始终低于 DeepEP v2。
- **完美均衡使其免受不平衡的影响**:随着 maxvio 的增长,MoonEP 的通信时间几乎保持平坦,而 DeepEP v2 —— 其延迟由最繁忙的 rank 决定 —— 则稳步下降。
- **该对比包含了 MoonEP 额外的 kernel**:MoonEP 增加了 DeepEP 不需要的规划和权重预取 kernel,并且它们已经堆叠在上面的柱状图中。即使算上整个关键路径,总的 dispatch 时间也与 DeepEP v2 单独的 dispatch 相当,并且在不平衡情况下表现更优,同时 combine 在各个级别上都显著更快。
**端到端训练**:
- **DeepEP 随着不平衡而退化**:最繁忙的 rank 接收更多的 token,因此迭代时间随着 maxvio 的增长而稳步攀升;与此同时,不断变化的激活 shape 会碎片化 GPU 显存,直到在高不平衡度下训练发生 OOM。
- **MoonEP 不受影响**:每个 rank 在每层始终精确计算 `S × K` 个 token,因此迭代时间在任何不平衡级别都保持平稳;完全静态的内存 shape 意味着没有碎片化,训练也永远不会 OOM。
## 支持的设备
- NVIDIA GPU
- 玄武 PPU(审核中,即将推出)
## 用法
### 集成
**符号说明**:`S` = 每个 rank 的输入 token 数,`K` = 每个 token 路由的 top-k,`E` = EP 组中路由专家总数,`R` = EP rank 数(EP 通信大小),`B` = 每个 rank 的权重预取槽位,`NvS` = 每个 rank 分发的 token 槽位(`S × K` 个真实 token 加上每个 VM 组的 padding),`H` = 隐藏层大小,`H'` = 专家 FFN 中间层大小。
MoonEP 与训练或推理框架的契约是 **每个专家投影对应一个连续的对称内存权重张量,外加一个由规划器生成的 `cu_seqlens`**。VM 组 GEMM 消耗一个单一的 `[E+B, H, H']` 权重张量;`cu_seqlens[E+B]`(由 `dispatch` 返回)选择当前步骤激活哪些专家行。
#### 权重 buffer
对于每个专家投影(gate/up/down),每一层都持有一个 **连续的 VMM 范围** `[E+B, H, H']`,并且在每个 rank 上的布局都是相同的。连续性是一个硬性要求:组 GEMM 完全通过行索引来定位专家。
- **行 `[0, E)`:所有 rank 的本地专家** —— 每个 rank `E/R` 行。每个数据块在物理上*就是*宿主 rank 的参数内存,并通过对称内存映射到各处。
- **行 `[E, E+B)`:本地预取槽位**,由 `buffer.prefetch_weight` 填充;规划器通过 `cu_seqlens` 将重复专家的 token 片段指向这些槽位。它们的物理内存来自一个所有层共享的进程全局池,因此额外成本是每个投影总共 `B` 个专家权重,而不是每层。
**如何设置 B。**
- **训练**:必须使用 **`B = E/R`** —— 规划器最多只从每个 rank 的一个远程宿主组复制专家(≤ `E/R` 个专家),因此组 GEMM 涉及的每个专家都是本地的。
- **推理**(仅预取,无梯度):允许 `B < E/R`,**推荐 `B = 3–4`**。如果某个 rank 需要的不同远程专家数量超过了 `B`,组 GEMM 会通过对称映射(内存语义,由 `cu_seqlens` 寻址)直接从宿主 rank 读取溢出的权重 —— 稍慢一些,但对正确性没有影响。
#### 梯度 buffer(仅限训练)
训练在 fp32 中镜像了权重布局:每个投影对应一个连续的 `[E+B, H, H']` **梯度 buffer**。
- **行 `[0, E)`**:所有者 rank 的参数梯度。
- **行 `[E, E+B)`**:预取槽位梯度,由 **一个单独的归约 buffer 支持,而不是由参数梯度支持** —— 重复专家的梯度是临时的,必须对框架自身的梯度归约保持不可见。与预取池一样,其物理内存来自一个跨层共享的进程全局池。
- **归约 buffer**:每个 rank 将所有 `R` 个归约 buffer 映射为一个 `[R, B, H, H']` 视图。`reduce_grad` 让每个 rank 从每个 rank 的归约 buffer 中读取持有其自身专家梯度的槽位(通过 NVLink 远程读取),将它们累加到本地参数梯度中,然后将自身已消耗的槽位清零,以供下一个 microbatch 使用。
### API 指南
```
from moonep import Buffer
buffer = Buffer(S=4096, H=7168, K=8, E=256, num_ep_ranks=8,
num_sms=32, token_padding=128)
```
- `num_sms=None` 默认为 32。`B` 默认为 `E // num_ep_ranks`;也可以传入明确的值,如 `B=4`。
- `dispatch` / `combine` / `prefetch_weight` / `reduce_grad` 都接受 `async_finish=True`,以便在通信流上运行并返回一个 CUDA event。
#### dispatch fwd
```
hidden_nvsh, route_weights_nvs, cu_seqlens, plan = buffer.dispatch(
hidden_sh, # [S, H] bf16
route_weights_sk, # [S, K] fp32
topk_experts_sk, # [S, K] int32
tokens_per_expert, # [E] int32, local count
)
# hidden_nvsh: [NvS, H] bf16 — 按 physical VM group 顺序 dispatch 的 tokens
# route_weights_nvs: [NvS] fp32
# cu_seqlens: [E+B] int32 — 每个 VM group 行的 padded token end offset
# plan: MoonEPCommPlan — 将其保存用于 prefetch/combine 和两次 backward pass
buffer.prefetch_weight(
plan=plan,
full_gate_weight=full_gate_weight, # [E+B, H, H'] bf16
full_up_weight=full_up_weight, # [E+B, H, H'] bf16
full_down_weight=full_down_weight, # [E+B, H, H'] bf16
)
# full_*_weight: 行 [0, E) 是 source expert weights,行 [E, E+B) 是 prefetch slots
```
#### dispatch bwd
dispatch 的反向:将每个 token 的 K 份分发梯度副本累加回以 token 为主的布局 —— 即一次 combine —— 并将重复专家的权重梯度归约回它们的宿主 rank。
```
grad_hidden_sh, _, _ = buffer.combine(
plan=plan,
hidden_nvsh=grad_hidden_nvsh, # [NvS, H] bf16
)
# grad_hidden_sh: [S, H] bf16
buffer.reduce_grad(
plan=plan,
full_gate_grad=full_gate_grad, # [E+B, H, H'] fp32
full_up_grad=full_up_grad, # [E+B, H, H'] fp32
full_down_grad=full_down_grad, # [E+B, H, H'] fp32
gate_reduce_buffer=gate_reduce_buffer, # [R, B, H, H'] fp32
up_reduce_buffer=up_reduce_buffer, # [R, B, H, H'] fp32
down_reduce_buffer=down_reduce_buffer, # [R, B, H, H'] fp32
)
# full_*_grad: 与 weights 相同的 [E+B] 布局;行 [E, E+B) 由 reduce buffer 支持
# *_reduce_buffer: 所有 R 个 rank 的 reduce buffer 映射为一个 [R, B, H, H'] view
```
#### combine fwd
```
output_sh, gathered_route_weights_sk, _ = buffer.combine(
plan=plan,
hidden_nvsh=expert_output_nvsh, # [NvS, H] bf16
route_weights_nvs=route_weights_nvs, # [NvS] fp32, optional
)
# output_sh: [S, H] bf16 — 合并后的 token-major output
# gathered_route_weights_sk: [S, K] fp32 或 None — 收集回 token-major 的 routing weights
```
#### combine bwd
combine 的反向:通过使用保存的计划重新 dispatch,将输出梯度散射回 VM 组的顺序 —— 跳过规划,也不需要预取。
```
grad_expert_output_nvsh, _, _, _ = buffer.dispatch(
grad_output_sh, # [S, H] bf16
plan=plan,
)
# grad_expert_output_nvsh: [NvS, H] bf16
```
#### zero_copy
默认情况下,`dispatch` 返回新的张量,而 `combine` 首先将其输入复制到 NVL 分片中。如果在两端都设置 `zero_copy=True`,`dispatch` 会返回通信 buffer 的视图,并且专家 FFN 会原地读写它们 —— 完全没有边界拷贝:
```
hidden_nvsh, route_weights_nvs, cu_seqlens, plan = buffer.dispatch(
hidden_sh, route_weights_sk, topk_experts_sk, tokens_per_expert,
zero_copy=True,
)
# hidden_nvsh / route_weights_nvs 是 NVL buffer 的 view;
# expert FFN 必须将其输出原地写入 hidden_nvsh
output_sh, gathered_route_weights_sk, _ = buffer.combine(
plan=plan,
hidden_nvsh=hidden_nvsh,
route_weights_nvs=route_weights_nvs,
zero_copy=True, # asserts the inputs are exactly the views from dispatch
)
```
- 这些视图别名了会被下一次 `dispatch` / `combine` 覆盖的 buffer 状态 —— 不要在通信调用之间持有它们(autograd 不得将它们保存用于反向传播;这种情况需要设置 `zero_copy=False`)。
```
# 在销毁 process group 之前,显式释放由 Buffer 持有的 VMM/NVLink 资源
buffer.destroy()
```
## 构建与测试
```
pip install -e .
# 运行测试(需要多个 GPU + NVLink)
torchrun --nproc_per_node=8 -m pytest tests/test_planning.py
torchrun --nproc_per_node=8 -m pytest tests/test_dispatch.py
torchrun --nproc_per_node=8 -m pytest tests/test_combine.py
torchrun --nproc_per_node=8 -m pytest tests/test_e2e.py
torchrun --nproc_per_node=8 -m pytest tests/test_grad_reduce.py
torchrun --nproc_per_node=8 -m pytest tests/test_prefetch.py
```
## 致谢
本库受到了以下作品的启发:
- [DeepEP](https://github.com/deepseek-ai/DeepEP)
- [Echo](https://arxiv.org/abs/2603.07685)
- [UltraEP](https://github.com/Dots-Infra/UltraEP)
## 引用
```
@misc{moonep2026,
title={MoonEP: A Perfectly Balanced Expert Parallelism Library via Dynamic Redundant Experts},
author={Yutian Chen, Cong Li, Yucheng Wang, Ming Wei},
year={2026},
publisher = {GitHub},
howpublished = {\url{https://github.com/MoonshotAI/MoonEP}},
}
```
- **零拷贝使原始通信更快**:token 被直接写入到它们在远程 rank 上最终的按专家分组的位置 —— 无需 permute 输入,也无需 permute 输出 —— 通信 buffer 的视图被直接交给计算过程,从而消除了在尾声阶段占主导地位的 通信 buffer → 用户 buffer 拷贝。MoonEP 的通信时间在每种不平衡级别下都始终低于 DeepEP v2。
- **完美均衡使其免受不平衡的影响**:随着 maxvio 的增长,MoonEP 的通信时间几乎保持平坦,而 DeepEP v2 —— 其延迟由最繁忙的 rank 决定 —— 则稳步下降。
- **该对比包含了 MoonEP 额外的 kernel**:MoonEP 增加了 DeepEP 不需要的规划和权重预取 kernel,并且它们已经堆叠在上面的柱状图中。即使算上整个关键路径,总的 dispatch 时间也与 DeepEP v2 单独的 dispatch 相当,并且在不平衡情况下表现更优,同时 combine 在各个级别上都显著更快。
**端到端训练**:
- **DeepEP 随着不平衡而退化**:最繁忙的 rank 接收更多的 token,因此迭代时间随着 maxvio 的增长而稳步攀升;与此同时,不断变化的激活 shape 会碎片化 GPU 显存,直到在高不平衡度下训练发生 OOM。
- **MoonEP 不受影响**:每个 rank 在每层始终精确计算 `S × K` 个 token,因此迭代时间在任何不平衡级别都保持平稳;完全静态的内存 shape 意味着没有碎片化,训练也永远不会 OOM。
## 支持的设备
- NVIDIA GPU
- 玄武 PPU(审核中,即将推出)
## 用法
### 集成
**符号说明**:`S` = 每个 rank 的输入 token 数,`K` = 每个 token 路由的 top-k,`E` = EP 组中路由专家总数,`R` = EP rank 数(EP 通信大小),`B` = 每个 rank 的权重预取槽位,`NvS` = 每个 rank 分发的 token 槽位(`S × K` 个真实 token 加上每个 VM 组的 padding),`H` = 隐藏层大小,`H'` = 专家 FFN 中间层大小。
MoonEP 与训练或推理框架的契约是 **每个专家投影对应一个连续的对称内存权重张量,外加一个由规划器生成的 `cu_seqlens`**。VM 组 GEMM 消耗一个单一的 `[E+B, H, H']` 权重张量;`cu_seqlens[E+B]`(由 `dispatch` 返回)选择当前步骤激活哪些专家行。
#### 权重 buffer
对于每个专家投影(gate/up/down),每一层都持有一个 **连续的 VMM 范围** `[E+B, H, H']`,并且在每个 rank 上的布局都是相同的。连续性是一个硬性要求:组 GEMM 完全通过行索引来定位专家。
- **行 `[0, E)`:所有 rank 的本地专家** —— 每个 rank `E/R` 行。每个数据块在物理上*就是*宿主 rank 的参数内存,并通过对称内存映射到各处。
- **行 `[E, E+B)`:本地预取槽位**,由 `buffer.prefetch_weight` 填充;规划器通过 `cu_seqlens` 将重复专家的 token 片段指向这些槽位。它们的物理内存来自一个所有层共享的进程全局池,因此额外成本是每个投影总共 `B` 个专家权重,而不是每层。
**如何设置 B。**
- **训练**:必须使用 **`B = E/R`** —— 规划器最多只从每个 rank 的一个远程宿主组复制专家(≤ `E/R` 个专家),因此组 GEMM 涉及的每个专家都是本地的。
- **推理**(仅预取,无梯度):允许 `B < E/R`,**推荐 `B = 3–4`**。如果某个 rank 需要的不同远程专家数量超过了 `B`,组 GEMM 会通过对称映射(内存语义,由 `cu_seqlens` 寻址)直接从宿主 rank 读取溢出的权重 —— 稍慢一些,但对正确性没有影响。
#### 梯度 buffer(仅限训练)
训练在 fp32 中镜像了权重布局:每个投影对应一个连续的 `[E+B, H, H']` **梯度 buffer**。
- **行 `[0, E)`**:所有者 rank 的参数梯度。
- **行 `[E, E+B)`**:预取槽位梯度,由 **一个单独的归约 buffer 支持,而不是由参数梯度支持** —— 重复专家的梯度是临时的,必须对框架自身的梯度归约保持不可见。与预取池一样,其物理内存来自一个跨层共享的进程全局池。
- **归约 buffer**:每个 rank 将所有 `R` 个归约 buffer 映射为一个 `[R, B, H, H']` 视图。`reduce_grad` 让每个 rank 从每个 rank 的归约 buffer 中读取持有其自身专家梯度的槽位(通过 NVLink 远程读取),将它们累加到本地参数梯度中,然后将自身已消耗的槽位清零,以供下一个 microbatch 使用。
### API 指南
```
from moonep import Buffer
buffer = Buffer(S=4096, H=7168, K=8, E=256, num_ep_ranks=8,
num_sms=32, token_padding=128)
```
- `num_sms=None` 默认为 32。`B` 默认为 `E // num_ep_ranks`;也可以传入明确的值,如 `B=4`。
- `dispatch` / `combine` / `prefetch_weight` / `reduce_grad` 都接受 `async_finish=True`,以便在通信流上运行并返回一个 CUDA event。
#### dispatch fwd
```
hidden_nvsh, route_weights_nvs, cu_seqlens, plan = buffer.dispatch(
hidden_sh, # [S, H] bf16
route_weights_sk, # [S, K] fp32
topk_experts_sk, # [S, K] int32
tokens_per_expert, # [E] int32, local count
)
# hidden_nvsh: [NvS, H] bf16 — 按 physical VM group 顺序 dispatch 的 tokens
# route_weights_nvs: [NvS] fp32
# cu_seqlens: [E+B] int32 — 每个 VM group 行的 padded token end offset
# plan: MoonEPCommPlan — 将其保存用于 prefetch/combine 和两次 backward pass
buffer.prefetch_weight(
plan=plan,
full_gate_weight=full_gate_weight, # [E+B, H, H'] bf16
full_up_weight=full_up_weight, # [E+B, H, H'] bf16
full_down_weight=full_down_weight, # [E+B, H, H'] bf16
)
# full_*_weight: 行 [0, E) 是 source expert weights,行 [E, E+B) 是 prefetch slots
```
#### dispatch bwd
dispatch 的反向:将每个 token 的 K 份分发梯度副本累加回以 token 为主的布局 —— 即一次 combine —— 并将重复专家的权重梯度归约回它们的宿主 rank。
```
grad_hidden_sh, _, _ = buffer.combine(
plan=plan,
hidden_nvsh=grad_hidden_nvsh, # [NvS, H] bf16
)
# grad_hidden_sh: [S, H] bf16
buffer.reduce_grad(
plan=plan,
full_gate_grad=full_gate_grad, # [E+B, H, H'] fp32
full_up_grad=full_up_grad, # [E+B, H, H'] fp32
full_down_grad=full_down_grad, # [E+B, H, H'] fp32
gate_reduce_buffer=gate_reduce_buffer, # [R, B, H, H'] fp32
up_reduce_buffer=up_reduce_buffer, # [R, B, H, H'] fp32
down_reduce_buffer=down_reduce_buffer, # [R, B, H, H'] fp32
)
# full_*_grad: 与 weights 相同的 [E+B] 布局;行 [E, E+B) 由 reduce buffer 支持
# *_reduce_buffer: 所有 R 个 rank 的 reduce buffer 映射为一个 [R, B, H, H'] view
```
#### combine fwd
```
output_sh, gathered_route_weights_sk, _ = buffer.combine(
plan=plan,
hidden_nvsh=expert_output_nvsh, # [NvS, H] bf16
route_weights_nvs=route_weights_nvs, # [NvS] fp32, optional
)
# output_sh: [S, H] bf16 — 合并后的 token-major output
# gathered_route_weights_sk: [S, K] fp32 或 None — 收集回 token-major 的 routing weights
```
#### combine bwd
combine 的反向:通过使用保存的计划重新 dispatch,将输出梯度散射回 VM 组的顺序 —— 跳过规划,也不需要预取。
```
grad_expert_output_nvsh, _, _, _ = buffer.dispatch(
grad_output_sh, # [S, H] bf16
plan=plan,
)
# grad_expert_output_nvsh: [NvS, H] bf16
```
#### zero_copy
默认情况下,`dispatch` 返回新的张量,而 `combine` 首先将其输入复制到 NVL 分片中。如果在两端都设置 `zero_copy=True`,`dispatch` 会返回通信 buffer 的视图,并且专家 FFN 会原地读写它们 —— 完全没有边界拷贝:
```
hidden_nvsh, route_weights_nvs, cu_seqlens, plan = buffer.dispatch(
hidden_sh, route_weights_sk, topk_experts_sk, tokens_per_expert,
zero_copy=True,
)
# hidden_nvsh / route_weights_nvs 是 NVL buffer 的 view;
# expert FFN 必须将其输出原地写入 hidden_nvsh
output_sh, gathered_route_weights_sk, _ = buffer.combine(
plan=plan,
hidden_nvsh=hidden_nvsh,
route_weights_nvs=route_weights_nvs,
zero_copy=True, # asserts the inputs are exactly the views from dispatch
)
```
- 这些视图别名了会被下一次 `dispatch` / `combine` 覆盖的 buffer 状态 —— 不要在通信调用之间持有它们(autograd 不得将它们保存用于反向传播;这种情况需要设置 `zero_copy=False`)。
```
# 在销毁 process group 之前,显式释放由 Buffer 持有的 VMM/NVLink 资源
buffer.destroy()
```
## 构建与测试
```
pip install -e .
# 运行测试(需要多个 GPU + NVLink)
torchrun --nproc_per_node=8 -m pytest tests/test_planning.py
torchrun --nproc_per_node=8 -m pytest tests/test_dispatch.py
torchrun --nproc_per_node=8 -m pytest tests/test_combine.py
torchrun --nproc_per_node=8 -m pytest tests/test_e2e.py
torchrun --nproc_per_node=8 -m pytest tests/test_grad_reduce.py
torchrun --nproc_per_node=8 -m pytest tests/test_prefetch.py
```
## 致谢
本库受到了以下作品的启发:
- [DeepEP](https://github.com/deepseek-ai/DeepEP)
- [Echo](https://arxiv.org/abs/2603.07685)
- [UltraEP](https://github.com/Dots-Infra/UltraEP)
## 引用
```
@misc{moonep2026,
title={MoonEP: A Perfectly Balanced Expert Parallelism Library via Dynamic Redundant Experts},
author={Yutian Chen, Cong Li, Yucheng Wang, Ming Wei},
year={2026},
publisher = {GitHub},
howpublished = {\url{https://github.com/MoonshotAI/MoonEP}},
}
```
标签:Vectored Exception Handling, 专家并行, 人工智能, 凭据扫描, 分布式计算, 大模型, 异常处理, 混合专家模型, 用户模式Hook绕过, 逆向工具, 通信优化