news 2026/9/29 6:24:18

深入浅出 DeepSeek MoE:EP 与 FSDP 经典二次开发实战指南

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
深入浅出 DeepSeek MoE:EP 与 FSDP 经典二次开发实战指南
  • 文档
  • 教程
  • 人工智能
  • 大模型
  • RLHF

【免费下载链接】Awesome-ML-SYS-Tutorial

My learning notes for ML SYS.

项目地址:https://gitcode.com/gh_mirrors/aw/Awesome-ML-SYS-Tutorial
点击查看免费下载

本指南以当前仓库 rlhf/sys-design/readme-4.md 为骨架,围绕"MoE 稀疏激活、专家并行(EP)的 All-to-All 通信原理、EP 与 TP 的量化对比,以及如何在仅支持 DP 的 FSDP 上二次开发 EP"这条主线展开。全文结合仓库中 slime FSDP 后端、FSDP2 原理 与 权重更新机制 等资源,深入剖析 VeOmni、Automodel、TorchTitan 三个开源社区项目的 EP+FSDP 实现,读者读完后能够系统理解 MoE 模型并行策略选型的核心权衡,并掌握在 FSDP 上叠加 EP 时涉及的关键代码模式(专家切分、All-to-All 调度、prefetch 配置、DeepEP 集成)。

写作背景:为什么要在 FSDP 上二次开发 EP

FSDP(Fully Sharded Data Parallel)本质上只支持 DeepSpeed ZeRO 类型的数据并行,TP、PP、EP 均无官方实现,需要在 HuggingFace Transformers 生态上自行二次开发。这一需求在 MoE 模型盛行的当下尤为迫切:原生 Transformers + FSDP 在 MoE 模型上存在明显的性能短板,一批又一批工程师试图通过二次开发来弥补。

本文记录的正是这一调研与实践过程:作者所在的 SGLang RL 小组(即 slime 框架的开发团队)曾多次讨论,是否要在 slime 已经支持的 FSDP 训练后端之上进一步支持 EP。作为背景对照,仓库中的 slime FSDP 后端文档 明确写到:FSDP 后端目前仅支持 DP + CP,不支持 TP、EP、PP,且"未来计划"中第一条就是"维持代码的干净与整洁的同时实现 TP 和 EP"——这正是本文讨论的二次开发动机。下文先讲清楚 MoE 与 EP 本身,再给出社区经典实现的学习笔记。

DeepSeek MoE:稀疏激活时代的开端

在 Dense 模型中,每一层的所有参数都会参与每个 Token 的计算。MoE 架构则将原本巨大的全连接层 FFN 拆分为多个规模较小、结构相同的独立单元——专家(Experts),并引入**稀疏激活(Sparse Activation)**机制:对于输入的每一个 Token,只有一小部分专家(如 Top-k)会被选中参与计算。这使得模型可以在保持计算量(FLOPs)基本不变的前提下,通过增加专家数量极大地扩张参数总量,某种意义上兼具更高的模型能力上限与更低的计算开销。

以 DeepSeek MoE(发表于 2024 年初,堪称 MoE 统治时代的开端之作)为代表的论文,在基础 MoE 之上进一步引入了两点关键优化:

  1. 共享专家(Shared Experts):在任意 forward 过程中,每层有少量专家永远被激活。某种意义上,这些 shared experts 存储着"常识"。
  2. 细粒度专家(Fine-grained Experts):相较于传统 MoE,将专家拆得更细。例如之前一层 FFN 拆成 8 个专家,现在拆成 64 个专家,在 DeepSeek V3 中甚至达到 256 个专家。

DeepSeek MoE 论文另一个令人印象深刻之处,是 MoE 与 Dense 模型之间的公平比较。不能拿 Llama 3.1 405B 与 DeepSeek V3 直接对比并宣称"MoE 强于 Dense",因为二者的变量差距远不止 MoE/Dense 一项。严格的控制变量必须从 pretrain 阶段开始,类似《Physics of Language Model》的做法。DeepSeek MoE 正是从 pretrain 的 token 数量开始做控制变量,得出结论:对于总参数量为 X、激活参数量为 Y 的 MoE 模型,其表现能够高于总参数量为 Y 的 Dense 模型,计算开销低于总参数量为 X 的 Dense 模型,甚至有接近并超越总参数量为 X 的 Dense 模型的可能性。

当然,如果 MoE 模型没有训好,部分 experts 在推理过程中一直不激活,就会出现下图的尴尬局面:

对于这种"僵尸专家",甚至可以直接剪枝掉从未被使用的 experts:总参数量下降,能力却不降低。

Expert Parallelism:把专家物理拆分到不同 Rank

先回顾朴素的 MoE forward 流程:

  1. Gate(路由)计算:输入 Token 经过 Gate 网络,计算该 Token 与各个专家的相关性得分。
  2. 专家选择:Gate 根据得分选出需要参与计算的专家。
  3. Token 分发:Token 被发送至选中的专家。
  4. 专家并行计算:选中的专家各自独立完成矩阵乘法运算。
  5. 结果合并:将各专家的输出按 Gate 权重加权求和,传入下一层。

在没有 EP 的情况下,每张 GPU 都必须存储该层所有专家的完整权重。对于参数量动辄上千亿、拥有数百个专家的 MoE 模型,单卡显存完全无法承载。因此EP 的核心逻辑是将 Experts 集合在第 0 维(专家维度)拆分,让不同 Rank 维护不同的专家子集。由于专家被物理隔离在不同显卡上,原本 gate 所做的逻辑分发变成了真正跨越 GPU rank 的物理分发。

下图展示了带设备部署的 MoE Transformer 编码器:MoE 层被切分到多个 Device 上,并通过All-to-All Dispatch与All-to-All Combine两段通信完成 Token 与专家之间的跨设备交互:

EP 的完整数据流分为三个阶段:

1. Dispatch(All-to-All)

在训练阶段,虽然 Transformer 是 auto-regressive 的,但 Causal Mask 实现了全序列的并行 Forward。此时每个 Rank 同时持有大量 Token,且每个 Token 持有不同专家的 gate 分数。每个 Rank 独立运行 Gate 算法,计算本地 Token 所需的目标专家及其所在的 Rank。

由于每个 Rank 都有 Token 需要发往其他 Rank,同时也要接收来自其他所有 Rank 的 Token,这构成了典型的All-to-All 通信,完成了**分布式转置(Distributed Transpose)**过程:将数据分布从"按序列位置对齐"重组为"按 experts 索引对齐"。

2. Expert Compute

各 Rank 并行执行本地持有的专家计算(FFN),不同 Rank 之间的专家计算完全独立,无需通信同步。但此阶段的计算效率极大取决于路由分布:如果大量 Token 涌向同一个专家(Hotspot Expert),会造成严重的负载不均衡(Load Imbalance),计算慢的 Rank 会拖累整个集群的同步速度。

3. Combine(All-to-All)

计算完成后,专家的局部输出(Expert Latents)需要再次通过 All-to-All 通信回传:每个专家 Rank 将计算结果发回给该 Token 原始出发的 Rank,确保数据物理分布恢复到进入 MoE 层之前的状态,以便进行后续的加权求和(Combine)、残差连接以及下一层 Attention 的并行计算。

注意到:虽然在 inference 的 Decoding 阶段每次仅处理一个 Token,但在 Training 和 Pre-filling 阶段,高吞吐的并行计算使 All-to-All 成为最有效、最常用的通信抽象。因此一般认为EP 需要两次 All-to-All 通信。

EP vs TP:为什么 MoE 更青睐 EP

TP 同样可以降低单 Rank 显存负载,甚至也能切分专家权重(将每个专家的参数分到不同 rank 上)。那为什么还需要 EP?二者的本质区别体现在通讯量与计算效率两个维度。

通讯量:TP 恒定 ≈2S,EP 随 k/N 缩放

定义以下变量:

  • $N$:并行组内的 GPU 数量(TP 组或 EP 组大小)。
  • $B \times L$:总 Token 数量(Batch Size × Sequence Length)。
  • $H$:Hidden Size(每个 Token 的向量维度)。
  • $k$:MoE 的 Top-$k$ 激活数(每个 Token 选择的专家数)。
  • $S = B \times L \times H$:该层输入数据的总激活量(Activation Size)。

基于主流的 Ring All-Reduce 与 Standard Exchange All-to-All 计算每个 GPU 发送的数据量:

TP将每个 expert 的 FFN 矩阵切分,每经过一个 expert 层,需要在 $W_{down}$ 之后做一次 All-Reduce。Ring 算法下单卡通讯量为 $2 \times \frac{N-1}{N} \times S$,即 $\text{Comm}_{TP} \approx 2S$。注意TP 的通讯量与 k(激活 expert 数)完全无关——哪怕只激活 1 个 expert,TP 也要雷打不动地同步全量激活值。

EP在专家维度显式拆分,Dispatch 阶段每个 Rank 最初持有的数据量仅为 $S/N$;每个 Token 需要分发给 k 个专家,单卡发出的数据量平均为 $k \times (S/N)$;算上 Combine 阶段的回传,单卡通讯量为:

$$\text{Comm}_{EP} = 2 \times \frac{N-1}{N} \times \frac{kS}{N} \approx \frac{2k}{N}S$$

在大规模并行时(如 $N=64$ 或 $256$),只要 $k < N$(DeepSeek 每层激活 $k=8$,专家总数 256),EP 的单卡通讯字节数远小于 TP。

通讯瓶颈:EP 的带宽与连接数劣势

虽然 EP 通讯量更少,但它遇到的通讯瓶颈更严重:

  1. 带宽差一个数量级:TP 通常死磕在单机 8 卡的 NVLink 域内,带宽起步 900GB/s;EP 往往要跨越节点走 RDMA,带宽通常只有 50GB/s ~ 100GB/s。
  2. All-to-All 是 $N^2$ 个小连接:在大规模集群下,握手开销、长尾延迟以及负载不均衡(Load Imbalance)带来的等待时间,可能比实际通讯耗时还要大。

计算效率:EP 保持算子完整,TP 产生"瘦长 GEMM"

TP 的核心逻辑是将矩阵"横着切"或"竖着切"。在 MoE 场景下,单个专家参数量通常较小;如果使用 TP,每个专家原本就不大的 $H \times \text{Hidden_Size}$ 矩阵会被进一步切成 $1/N$,在 GPU 上执行的是极其"瘦长"的矩阵乘法(GEMM)。对 NVIDIA Tensor Cores 而言,过小的维度无法充分填充计算流水线,实际算力利用率大幅下降。这也是 TP 一般不做跨机的一大原因:跨机通讯远比机内慢,对 TP 几乎恒定的通讯量而言是灾难;且 TP 切到更多机器会让每个 rank 的形状更加瘦长,GEMM 效率进一步下降。

而EP 能保持 expert 矩阵的完整性:虽然 Token 需要跨卡搬运,但一旦到达目标 GPU,面对的是形状规整、足以触发高效计算内核的完整矩阵。对于 DeepSeek 这种**细粒度专家(Fine-grained Experts)**设计,每个专家极其微小,若再套用 TP,计算效率将退化到难以忍受的地步。

此外,DeepSeek MoE 选择 EP 还有更多的 infra 创新:shared experts 被隔离、不参与 EP 通信;DeepEP 为大规模 k 带来极致通讯隐藏:

  1. 计算与通讯极致重叠(Stream-K):当 Dispatch 的第一批 Token 到达 GPU 时,计算内核立即启动,而不是等待 16 个专家的所有数据全部到齐。
  2. RDMA 直驱:DeepEP 绕过传统 NCCL 协议栈,利用 PTX 级别优化实现低延迟跨节点数据交换。在 $k=16$ 产生的巨大吞吐下,DeepEP 依然能维持极高的带宽利用率,使"通讯时间"几乎完全被"计算时间"掩盖。

综合对比如下:

维度TP 方案EP 方案 (DeepSeek 为例)
通讯量 (Bytes)固定 (≈2S)随 k/N 缩放 (≈2kS/N)
通讯延迟 (Latency)极高(高频 All-Reduce 同步锁)可控(粗粒度 All-to-All 异步掩盖)
计算效率 (MFU)低(矩阵切分导致算子不饱满)高(算子完整,易于硬件加速)
集群扩展性局限于单机 NVLink 域支持万卡集群 RDMA 扩展

需要强调,EP 优于 TP 的结论很大程度上是DeepSeek MoE 引导的"小而多"设计驱动的:单个 expert 小,TP 切细后 GEMM 效率暴跌;而 GShard 时代"大而少"的 MoE 单个 expert 大,TP 切分后 GEMM 效率仍有保证,因此当时对 MoE 采用 TP 才是主流。说到底,EP 是被算法驱动的新并行策略。

ETP:先 EP 再对每个专家做 TP

EP 将不同 experts 物理分到不同 rank 后,这些 experts 仍旧可以做 TP。例如 EP2-TP4:先把所有专家分成两组(0~3 卡负担前 1/2 的专家),单个专家再拆分到 4 个组内 rank 上。这种策略有专门术语ETP(先做 EP,再对每个 expert 做 TP)。

但 ETP 并没有在开源社区被广泛采用。以 2026 年 1 月 1 日的 SGLang cookbook 中启动 DeepSeek R1 的指令为例,TP 与 EP 可以同时开启:

python3 -m sglang.launch_server \ --model-path deepseek-ai/DeepSeek-R1-0528 \ --tp 8 \ --ep 8 \ --enable-symm-mem # Optional: improves performance, but may be unstable

这条指令中,TP 8 与 EP 8 的含义是:experts 被分到 8 个 rank 上(EP),而非 MoE 部分(如 linear 层)按 TP 分到 8 个 rank 上,执行的并不是 ETP 模式。理解这一点,是读懂现代推理引擎并行配置组合的关键。

经典的 FSDP 二次开发:EP

回到本文的出发点。有了 EP 后,MoE 层的 forward 与 backward 流程对比如下:

forwardbackward
gate

All-to-All Dispatch

expert compute (FSDP2)

all gather

Expert FFN compute



release

All-to-All Return

merge
gate

All-to-All Combine

expert compute (FSDP2)

all gather

Expert FFN compute

reduce-scatter

release

All-to-All Return

merge

forward 与 backward 没有显著区别,核心差异就是加入了 EP 的 all-to-all 通讯。注意两点:

  • Backward 中的 Reduce-Scatter 是 FSDP 的标准动作,与 EP 无关:计算出的完整梯度需要在 DP 组内聚合(Reduce)并重新切分(Scatter)回各个 Rank,在数学上完成梯度的平均与分发。在 MoE EP 场景下,只有专家被进一步 FSDP 切分时,才有对应的 FSDP 级别 Reduce-Scatter。
  • experts 本身不需要在 EP 组内做梯度聚合:各自优化各自的梯度即可。

隐式切分 vs 显式切分:FSDP 与 EP 的本质差异

理解 EP 二次开发,必须先理解 FSDP 与 EP 切分方式的根本差异,这一点在仓库 FSDP 训练后端 中有充分铺垫:FSDP2 将每个参数表示为独立的DTensor,并在第 0 维上进行分片,保留了原始张量的全部元数据(shape、stride、dtype、placement 等)。

  • FSDP 的fully_shard是隐式切分(动态逻辑切分),希望对上层的模型代码无感。以shape=[128, H, I]的一组专家为例,在模型 forward 开始前,FSDP 内部会偷偷发起一次 all-gather,临时把 8 张卡上的碎片拼回完整的[128, H, I];对于上层而言,看到的还是一个完整 Tensor,无需关心分布式通信。
  • EP 是显式切分(静态物理切分):原本shape=[128, H, I]的一组专家,在每个 Rank 上物理变成[32, H, I]。模型代码必须感知到这个变化——MoE 层代码一定要知道"我这台机器上只有 32 个专家",并据此计算。

社区通用的五大优化方向

遍览各大框架,基于 FSDP 的 EP 二次开发通常包含以下优化:

  1. EP 切分 dim 0 experts,FSDP 切分 dim 1 hidden size(见 VeOmni):EP 按专家维切分,FSDP 再按隐藏维切分专家权重,两级切分正交叠加。
  2. Prefetch:在计算第 n 层时预先把第 n+1 层参数 gather 起来。只用 FSDP 可以直白地做 prefetch;但因为 EP 计算开始前有通信,需要手动操作保证前向和反向的 prefetch(见 VeOmni)。
  3. DeepEP:苦 NCCL 久矣,使用 RDMA 直驱替代 NCCL All-to-All(见 Automodel)。
  4. EPLB:通过专家冗余解决专家计算负载不均衡。
  5. Fused MoE:与 EP 关系不大,单个 GPU 负责多个专家时用 Fused MoE kernel 加速这些专家的计算。

实现对比:三个社区高光项目

以下对比 VeOmni、TorchTitan、Automodel 三个社区项目,它们都是先 EP,然后对 EP 完的每个块做 FSDP,是同一设计范式下的三种不同工程侧重。

VeOmni:EP + FSDP2 的整合样板

VeOmni 的代码结构如下:

VeOmni/veomni/ ├── distributed/ │ ├── parallel_state.py ← 全局并行状态(ep_fsdp_device_mesh, ep_size) │ ├── parallel_plan.py ← EP 切分计划(ParallelPlan.apply()) │ ├── torch_parallelize.py ← EP + FSDP 整合入口 │ │ ├── parallelize_model_fsdp2() ← 主入口 │ │ └── 手动 prefetch 配置 │ ├── fsdp/ │ │ ├── clip_grad_norm.py ← FSDP1 EP 感知梯度裁剪 │ │ └── extension.py ← Checkpoint 扩展 │ └── fsdp2/ │ └── clip_grad_norm.py ← FSDP2 EP 感知梯度裁剪 ├── models/ │ └── transformers/ │ └── qwen3_moe/ │ └── parallel_plan.py ← 模型特定的 EP 参数定义 └── sequence_parallel/ ├── async_ulysses.py ← 异步序列并行(与 EP 无关) └── ulysses.py ← 标准 Ulysses

整体逻辑清晰:专家先 apply EP 在第 0 维(expert)切分,再 FSDP;非专家部分直接按常规 FSDP 即可。其核心注释概括了完整流程:

Applies EP (when enabled) + FSDP2 parallel strategy to the model. Flow: 1. Apply EP: Expert tensors [128,H,I] -> [32,H,I] local tensors per EP rank 2. Apply FSDP2 to expert modules: Shard expert tensors along dim-1 (hidden dim) 3. Apply FSDP2 to regular modules: Standard dim-0 sharding 4. Result: Expert params [32, H/fsdp_size, I], regular params use standard FSDP2

关键函数parallelize_model_fsdp2节选:

def parallelize_model_fsdp2(model, enable_mixed_precision=True, basic_modules=None, **kwargs): # 【1】专家 128 -> 32 (EP) if parallel_state.ep_enabled: parallel_plan = model.get_parallel_plan() parallel_plan.apply(model, parallel_state.ep_fsdp_device_mesh) experts_map = parallel_plan.get_fsdp_no_shard_info(model) # 【2. 循环分片】由内而外切分每一层 layer_pairs = [] for layer_fqn, layer_mod in decoder_blocks: experts_mod = next((exp_mod for exp_fqn, exp_mod in experts_map.items() if ...), None) layer_mod._fsdp_modules = [] if experts_mod: fully_shard(experts_mod, **expert_fsdp_kwargs) # 切专家 layer_mod._fsdp_modules.append(experts_mod) fully_shard(layer_mod, **fsdp_kwargs) # 切整层 layer_mod._fsdp_modules.append(layer_mod) layer_pairs.append(layer_mod) # 【3. 切root model】 fully_shard(model, **fsdp_kwargs) # 【4. 配置prefetch】 # 正向 for cur, nxt in zip(layer_pairs, layer_pairs[1:] + [None]): if nxt: cur.set_modules_to_forward_prefetch(list(reversed(nxt._fsdp_modules))) # 反向 rev_blocks = list(reversed(layer_pairs)) for cur, prev in zip(rev_blocks, rev_blocks[1:] + [None]): if prev: cur.set_modules_to_backward_prefetch(list(reversed(prev._fsdp_modules))) return model

这段代码体现了"由内而外"的切分顺序:先fully_shard专家模块,再fully_shard整层,最后切 root model;随后通过set_modules_to_forward_prefetch与set_modules_to_backward_prefetch手动配置前后向的预取链路。这种基于_fsdp_modules列表的手动 prefetch 配置,正是前文所说"EP 计算开始前有通信,需要手动操作"的落点。

接着是通讯逻辑。All-to-All 通讯的复杂度不低,VeOmni 分为三步:

  1. Preprocess:在传输重数据之前,通过 all_gather 交换元数据,计算出 Input Splits 和 Output Splits;
  2. Dispatch:根据路由索引在本地 Permute,利用dist.all_to_all完成传输,收到数据后再次 Sort;
  3. Combine:计算完成后,执行逆向的通信和 Unpermute 操作,将 Token 还原回原始序列顺序。
def preprocess(expert_mask, num_experts, ep_group): # expert_mask: [Batch, Tokens, Num_Experts] (哪些 token 去哪些专家) # 1. 算出本地要发给每个 rank 的 token 数量 (Input Splits) ep_size = ep_group.size() num_local_tokens_per_expert = expert_mask.sum(dim=(1, 2)) input_splits = num_local_tokens_per_expert.reshape(ep_size, -1).sum(dim=1).tolist() # 2. dist.all_gather: 收集所有卡上的 num_local_tokens_per_expert num_global_tokens_per_expert = torch.zeros(...) dist.all_gather_into_tensor(num_global_tokens_per_expert, num_local_tokens_per_expert, group=ep_group) # 3. 算出本地将从每个 rank 接收多少 token (Output Splits) rank = dist.get_rank(ep_group) my_experts_range = slice(rank * num_local_experts, (rank + 1) * num_local_experts) tokens_sent_to_me = num_global_tokens_per_expert[:, my_experts_range] output_splits = tokens_sent_to_me.sum(dim=1).tolist() return input_splits, output_splits, tokens_sent_to_me def token_pre_all2all(hidden_states, expert_mask, input_splits, output_splits, ...): # 1. 本地重排 (Permute) # local_permuted: [Token1_to_Exp1, Token2_to_Exp1, ..., TokenN_to_Exp99] local_permuted, _ = permute(hidden_states, expert_mask.sum(dim=1)) # 2. All-to-All # 发送:input_splits, 接收:output_splits global_permuted = all_to_all(ep_group, local_permuted, output_splits, input_splits) # 3. Sort by Expert global_permuted = sort_chunks_by_idxs(global_permuted, ...) return global_permuted # 准备好喂给 Group GEMM 了 def tokens_post_all2all(expert_outputs, input_splits, output_splits, ...): # 1. 算完的数据是按 Expert 排列的,要发回去得按来源 Rank 重排 expert_outputs = sort_chunks_by_idxs(expert_outputs, ...) # 2. All-to-All Return unpermute_outputs = all_to_all(ep_group, expert_outputs, input_splits, output_splits) # 3. Unpermute final_output = unpermute(unpermute_outputs, ...) return final_output

注意preprocess的精妙之处:num_local_tokens_per_expert在本地算好后,通过all_gather_into_tensor收集全局路由分布,再按本 rank 负责的专家区间my_experts_range反推出 Output Splits——通信元数据本身也是分布式计算的产物。

Automodel:DeepEP 的深度集成样板

Automodel 的代码结构如下:

Automodel/nemo_automodel/ ├── components/ │ ├── distributed/ │ │ └── fsdp2.py ← FSDP2Manager(moe_mesh 定义) │ └── moe/ │ ├── parallelizer.py ← EP + FSDP 整合入口 │ │ ├── ExpertParallel ← EP 类定义 │ │ ├── apply_ep() ← EP 切分 │ │ ├── apply_fsdp() ← FSDP 切分 │ │ └── parallelize_model() ← 主入口 │ ├── layers.py ← MoE 层实现 │ ├── fsdp_mixin.py ← MoE FSDP 同步 Mixin(PP 相关) │ └── megatron/ │ ├── token_dispatcher.py ← Token 调度(_DeepepManager) │ ├── fused_a2a.py ← DeepEP 封装(FusedDispatch/Combine) │ └── moe_utils.py ← permute/unpermute 工具

Automodel 对 DeepEP 的使用可圈可点:通过_DeepepManager集成 DeepEP,用 Fused Dispatch/Combine 算子替代 NCCL All-to-All,调用链为token_dispatcher.py -> MoEFlexTokenDispatcher -> _DeepepManager -> fused_dispatch。

_DeepepManager是有状态的通信上下文管理器,封装 DeepEP 库与上层模型逻辑之间的交互。在 dispatch 阶段,DeepEP 底层返回一个handle对象(包含通信布局信息);在 combine 阶段,直接取出self.handle传给底层:

class _DeepepManager(_DispatchManager): """DeepEP backend for token dispatch/combine""" def __init__(self, group, router_topk, num_experts, num_local_experts, ...): self.group = group self.num_experts = num_experts self.num_local_experts = num_local_experts # 本 EP 组的 expert 数量 if fused_dispatch is None: raise ImportError("DeepEP is not installed.") def setup_metadata(self, num_local_tokens, probs): """处理 routing map""" probs = probs.reshape(num_local_tokens, self.num_experts) self.token_probs, self.token_indices = torch.topk(probs, self.router_topk, dim=-1) def dispatch(self, hidden_states, async_finish=False, allocate_on_comm_stream=False): """Dispatch tokens to experts""" # DeepEP 要求 float32 self.token_probs = self.token_probs.float() # 调用 DeepEP 的 fused_dispatch (hidden_states, dispatched_indices, dispatched_probs, num_tokens_per_expert, handle) = fused_dispatch( hidden_states, self.token_indices, self.token_probs, self.num_experts, self.group, async_finish=async_finish, ) self.handle = handle # 保存用于 combine return hidden_states def combine(self, hidden_states, async_finish=False, allocate_on_comm_stream=False): """Combine expert outputs""" hidden_states, _ = fused_combine( hidden_states, self.group, self.handle, # 使用 dispatch 时保存的 handle async_finish=async_finish, ) self.handle = None return hidden_states

FusedDispatch是torch.autograd.Function的封装,在 forward 中获取 DeepEP Buffer、计算 dispatch layout、调用核心 dispatch 并保存 handle;在 backward 中调用buffer.combine完成梯度的回传——DeepEP 的 combine 即反向传播的 dispatch 逆操作:

class FusedDispatch(torch.autograd.Function): @staticmethod def forward(ctx, x, token_indices, token_probs, num_experts, group, async_finish, ...): # 获取 DeepEP Buffer buffer = get_buffer(group, get_hidden_bytes(x)) # 计算 dispatch layout (num_tokens_per_rank, num_tokens_per_rdma_rank, num_tokens_per_expert, is_token_in_rank, event) = buffer.get_dispatch_layout( token_indices, num_experts, ... ) # 调用 DeepEP 核心 dispatch (recv_x, recv_token_indices, recv_token_probs, num_recv_tokens_per_expert_list, handle, after_event) = buffer.dispatch( x, topk_idx=token_indices, topk_weights=token_probs, # 必须 float32 num_tokens_per_rank=num_tokens_per_rank, ... async_finish=async_finish, ) # 异步同步 if async_finish: after_event.current_stream_wait() ctx.handle = handle # 保存用于 backward return recv_x, recv_token_indices, recv_token_probs, tokens_per_expert, handle @staticmethod def backward(ctx, ...): # backward 调用 combine grad_x, grad_token_probs, after_event = buffer.combine( grad_output.contiguous(), ctx.handle, ... ) return grad_x, ...

值得注意的实现细节:DeepEP 的topk_weights必须为 float32(self.token_probs = self.token_probs.float()),这是 DeepEP 接口的硬性约束;async_finish与after_event.current_stream_wait()则是实现"计算掩盖通信"(Dispatch 第一批 Token 到达即启动计算)的关键机制。

TorchTitan:全链路 prefetch 的极致样板

torchtitan/ ├── distributed/ │ ├── expert_parallel.py ← 核心!EP 类定义 │ ├── parallel_dims.py ← Device Mesh 管理 │ └── deepep.py ← DeepEP 封装(可选) ├── models/ │ ├── moe/ │ │ ├── moe.py ← MoE 层实现 │ │ └── moe_deepep.py ← DeepEP MoE 变体 │ └── llama4/infra/ │ └── parallelize.py ← EP + FSDP 整合入口

TorchTitan 实现了全链路的 prefetch:

  1. MoE 感知预取:大多框架可能只预取下一层 Block,但 TorchTitan 在前向传播时会显式地同时预取下一层 Block 及其内部的 Experts([next_transformer_block, next_transformer_block.moe.experts]);
  2. 最大限度计算覆盖:从 Embedding 层到 Output 层,甚至在反向传播中都有对应的set_modules_to_backward_prefetch逻辑,用密不透风的预取最大限度让计算掩盖通信。

Parallelize 逻辑——注意专家 FSDP 切分时对shard_placement_fn的巧妙处理:当efsdp size × ep_degree超过专家总数时,专家权重无法沿 dim-0 继续切分,此时自动退化为沿 dim-1(hidden)切分:

for layer_id, transformer_block in model.layers.items(): if transformer_block.moe_enabled and ep_degree > 1: fsdp_mod_ep_config = fsdp_config.copy() fsdp_mod_ep_config["mesh"] = edp_mesh _experts_shard_placement_fn = None assert edp_mesh is not None assert hasattr(transformer_block, "moe") if ( edp_mesh["efsdp"].size() * ep_degree > transformer_block.moe.experts.num_experts ): _experts_shard_placement_fn = lambda param: Shard(1) fully_shard( transformer_block.moe.experts, **fsdp_mod_ep_config, reshard_after_forward=reshard_after_forward, shard_placement_fn=_experts_shard_placement_fn, ) transformer_block.moe.experts.set_gradient_divide_factor( gradient_divide_factor, ) fully_shard( transformer_block, **fsdp_config, reshard_after_forward=reshard_after_forward, )

Prefetch 逻辑——前向从 Embedding 出发逐层预取,MoE 层预取[next_block, next_block.moe.experts],最后一层预取[norm, output];反向则完全对称地从 Output 逆推回 Embedding:

transformer_blocks = list(model.layers.values()) next_transformer_blocks = transformer_blocks[1:] + [None] if model.tok_embeddings is not None and len(model.layers) > 0: model.tok_embeddings.set_modules_to_forward_prefetch([transformer_blocks[0]]) for transformer_block, next_transformer_block in zip( transformer_blocks, next_transformer_blocks ): if next_transformer_block is not None: if next_transformer_block.moe_enabled: transformer_block.set_modules_to_forward_prefetch( [next_transformer_block, next_transformer_block.moe.experts] ) else: transformer_block.set_modules_to_forward_prefetch( [next_transformer_block] ) elif model.norm is not None and model.output is not None: transformer_block.set_modules_to_forward_prefetch( [model.norm, model.output] ) # backward reversed_transformer_blocks = list(reversed(model.layers.values())) prev_transformer_blocks = reversed_transformer_blocks[1:] + [None] if model.norm is not None and model.output is not None and len(model.layers) > 0: model.output.set_modules_to_backward_prefetch([reversed_transformer_blocks[0]]) for transformer_block, prev_transformer_block in zip( reversed_transformer_blocks, prev_transformer_blocks ): if prev_transformer_block is not None: if prev_transformer_block.moe_enabled: transformer_block.set_modules_to_backward_prefetch( [prev_transformer_block, prev_transformer_block.moe.experts] ) else: transformer_block.set_modules_to_backward_prefetch( [prev_transformer_block] ) elif model.tok_embeddings is not None: transformer_block.set_modules_to_backward_prefetch([model.tok_embeddings])

从 EP 二次开发到完整 RL 系统

EP 二次开发只是 MoE 时代 RL 系统训练后端的一环。以当前仓库为参照,可串联起完整的知识链路:

  • FSDP2 原理基础:EP 叠加 FSDP 依赖DTensor、fully_shard的隐式切分与手动 prefetch 等机制,详见 RL 系统深思:FSDP 训练后端。
  • 实际落地的约束:slime 的 FSDP 后端目前仅支持 DP + CP,TP/EP/PP 仍在其未来计划中;FSDP 通过AutoModelForCausalLM.from_pretrained()自动读取 HuggingFace 架构信息,无需权重格式转换,这降低了二次开发的适配成本,见 slime FSDP 后端。
  • 权重同步闭环:EP 专家梯度各自优化后,最终仍要通过update_weights_from_tensor(handle tuple 序列化、跨进程传递、SGLang 侧重建 tensor)或分桶异步更新等方式同步回推理引擎,见 RL 系统深思:权重更新机制。
  • Disaggregated 场景的 EP 处理:在训练推理分离架构下,参数按 TP/PP/EP 三维并行交叉切分,训练端需构造与推理端一致的 Engine Replica 并行配置,并对 MoE 专家参数单独进行 EP AllGather 后再映射写入,见 RL 系统深思(权重传输篇)。

总结

本文完整梳理了从 DeepSeek MoE 稀疏激活架构到 FSDP 上 EP 二次开发的完整链路:MoE 通过"小而多"的细粒度专家实现算力与参数解耦;EP 通过两次 All-to-All 完成 Token 的分布式转置与回传,相比 TP 在通讯量上随 k/N 缩放、在算子效率上保持 GEMM 完整,是细粒度 MoE 时代的主流并行策略;而在 FSDP 上叠加 EP 属于显式切分与隐式切分的叠加,需要模型代码感知专家子集变化,并配套处理 All-to-All 调度、prefetch 与 DeepEP 集成。VeOmni 展示了 EP+FSDP2 的干净切分流程与手写 prefetch,Automodel 展示了 DeepEP 的handle状态管理与 autograd 封装,TorchTitan 则将前向/反向全链路 prefetch 做到极致——三者共同构成了一幅"在 FSDP 上二次开发 EP"的完整工程地图。

  • 文档
  • 教程
  • 人工智能
  • 大模型
  • RLHF

【免费下载链接】Awesome-ML-SYS-Tutorial

My learning notes for ML SYS.

项目地址:https://gitcode.com/gh_mirrors/aw/Awesome-ML-SYS-Tutorial
点击查看免费下载
上一篇:构建多语言表单:jQuery Validation国际化messages文件完全指南
下一篇:大模型面试笔记:从开源协作到知识共享的思考

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

版权声明: 本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若内容造成侵权/违法违规/事实不符,请联系邮箱:809451989@qq.com进行投诉反馈,一经查实,立即删除!
网站建设 2026/9/29 6:23:27

深度学习优化器全解析:从SGD到AdamW的选型与调参实战指南

1. Model-Optimizer到底在优化什么&#xff1a;先聊聊背景和真实痛点做深度学习训练的人应该都有过这种体验&#xff1a;模型结构照搬经典论文、数据也都干净&#xff0c;结果一跑起来loss死活不降&#xff0c;或者降一半突然NaN&#xff0c;再或者训练集都能满分、验证集却一路…

作者头像 李华
网站建设 2026/9/29 6:22:23

Agentic AI能跑Demo,为什么一上项目就崩?先把这三笔账算清楚

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

作者头像 李华
网站建设 2026/9/29 6:21:24

Hatch 构建配置完全指南:从文件选择到可复现构建

开发工具构建工具 【免费下载链接】hatch Modern, extensible Python project management 项目地址&#xff1a; https://gitcode.com/gh_mirrors/ha/hatch 点击查看 免费下载 本篇技术指南以 Hatch 项目的 docs/config/build.md 为骨架&#xff0c;系统讲解构建配置的核心主题…

作者头像 李华
网站建设 2026/9/29 6:19:22

TensorFlow工业落地实战:从环境配置到边缘部署全链路避坑指南

1. 这不是“又一个深度学习框架”——TensorFlow 是怎么从实验室走向产线的你搜“tensorflow”&#xff0c;页面上跳出来的全是安装报错、版本冲突、CUDA不匹配、GPU识别失败……但真正用过三年以上 TensorFlow 的人&#xff0c;第一反应不是“怎么装”&#xff0c;而是“这个模…

作者头像 李华