ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

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

深入浅出 DeepSeek MoE:EP 与 FSDP 经典二次开发实战指南 文档教程人工智能大模型RLHF【免费下载链接】Awesome-ML-SYS-TutorialMy 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 三个开源社区项目的 EPFSDP 实现读者读完后能够系统理解 MoE 模型并行策略选型的核心权衡并掌握在 FSDP 上叠加 EP 时涉及的关键代码模式专家切分、All-to-All 调度、prefetch 配置、DeepEP 集成。写作背景为什么要在 FSDP 上二次开发 EPFSDPFully 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 之上进一步引入了两点关键优化共享专家Shared Experts在任意 forward 过程中每层有少量专家永远被激活。某种意义上这些 shared experts 存储着常识。细粒度专家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 流程Gate路由计算输入 Token 经过 Gate 网络计算该 Token 与各个专家的相关性得分。专家选择Gate 根据得分选出需要参与计算的专家。Token 分发Token 被发送至选中的专家。专家并行计算选中的专家各自独立完成矩阵乘法运算。结果合并将各专家的输出按 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. DispatchAll-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. CombineAll-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 更青睐 EPTP 同样可以降低单 Rank 显存负载甚至也能切分专家权重将每个专家的参数分到不同 rank 上。那为什么还需要 EP二者的本质区别体现在通讯量与计算效率两个维度。通讯量TP 恒定 ≈2SEP 随 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 个 expertTP 也要雷打不动地同步全量激活值。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$$在大规模并行时如 $N64$ 或 $256$只要 $k N$DeepSeek 每层激活 $k8$专家总数 256EP 的单卡通讯字节数远小于 TP。通讯瓶颈EP 的带宽与连接数劣势虽然 EP 通讯量更少但它遇到的通讯瓶颈更严重带宽差一个数量级TP 通常死磕在单机 8 卡的 NVLink 域内带宽起步 900GB/sEP 往往要跨越节点走 RDMA带宽通常只有 50GB/s ~ 100GB/s。All-to-All 是 $N^2$ 个小连接在大规模集群下握手开销、长尾延迟以及负载不均衡Load Imbalance带来的等待时间可能比实际通讯耗时还要大。计算效率EP 保持算子完整TP 产生瘦长 GEMMTP 的核心逻辑是将矩阵横着切或竖着切。在 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 带来极致通讯隐藏计算与通讯极致重叠Stream-K当 Dispatch 的第一批 Token 到达 GPU 时计算内核立即启动而不是等待 16 个专家的所有数据全部到齐。RDMA 直驱DeepEP 绕过传统 NCCL 协议栈利用 PTX 级别优化实现低延迟跨节点数据交换。在 $k16$ 产生的巨大吞吐下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 再对每个专家做 TPEP 将不同 experts 物理分到不同 rank 后这些 experts 仍旧可以做 TP。例如 EP2-TP4先把所有专家分成两组03 卡负担前 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 流程对比如下forwardbackwardgateAll-to-All Dispatchexpert compute (FSDP2)all gatherExpert FFN computereleaseAll-to-All ReturnmergegateAll-to-All Combineexpert compute (FSDP2)all gatherExpert FFN computereduce-scatterreleaseAll-to-All Returnmergeforward 与 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 二次开发通常包含以下优化EP 切分 dim 0 expertsFSDP 切分 dim 1 hidden size见 VeOmniEP 按专家维切分FSDP 再按隐藏维切分专家权重两级切分正交叠加。Prefetch在计算第 n 层时预先把第 n1 层参数 gather 起来。只用 FSDP 可以直白地做 prefetch但因为 EP 计算开始前有通信需要手动操作保证前向和反向的 prefetch见 VeOmni。DeepEP苦 NCCL 久矣使用 RDMA 直驱替代 NCCL All-to-All见 Automodel。EPLB通过专家冗余解决专家计算负载不均衡。Fused MoE与 EP 关系不大单个 GPU 负责多个专家时用 Fused MoE kernel 加速这些专家的计算。实现对比三个社区高光项目以下对比 VeOmni、TorchTitan、Automodel 三个社区项目它们都是先 EP然后对 EP 完的每个块做 FSDP是同一设计范式下的三种不同工程侧重。VeOmniEP 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_precisionTrue, basic_modulesNone, **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 分为三步Preprocess在传输重数据之前通过 all_gather 交换元数据计算出 Input Splits 和 Output SplitsDispatch根据路由索引在本地 Permute利用dist.all_to_all完成传输收到数据后再次 SortCombine计算完成后执行逆向的通信和 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(dim1).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, groupep_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(dim1).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(dim1)) # 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——通信元数据本身也是分布式计算的产物。AutomodelDeepEP 的深度集成样板Automodel 的代码结构如下Automodel/nemo_automodel/ ├── components/ │ ├── distributed/ │ │ └── fsdp2.py ← FSDP2Managermoe_mesh 定义 │ └── moe/ │ ├── parallelizer.py ← EP FSDP 整合入口 │ │ ├── ExpertParallel ← EP 类定义 │ │ ├── apply_ep() ← EP 切分 │ │ ├── apply_fsdp() ← FSDP 切分 │ │ └── parallelize_model() ← 主入口 │ ├── layers.py ← MoE 层实现 │ ├── fsdp_mixin.py ← MoE FSDP 同步 MixinPP 相关 │ └── 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_finishFalse, allocate_on_comm_streamFalse): 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_finishasync_finish, ) self.handle handle # 保存用于 combine return hidden_states def combine(self, hidden_states, async_finishFalse, allocate_on_comm_streamFalse): Combine expert outputs hidden_states, _ fused_combine( hidden_states, self.group, self.handle, # 使用 dispatch 时保存的 handle async_finishasync_finish, ) self.handle None return hidden_statesFusedDispatch是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_idxtoken_indices, topk_weightstoken_probs, # 必须 float32 num_tokens_per_ranknum_tokens_per_rank, ... async_finishasync_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必须为 float32self.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 实现了全链路的 prefetchMoE 感知预取大多框架可能只预取下一层 Block但 TorchTitan 在前向传播时会显式地同时预取下一层 Block 及其内部的 Experts[next_transformer_block, next_transformer_block.moe.experts]最大限度计算覆盖从 Embedding 层到 Output 层甚至在反向传播中都有对应的set_modules_to_backward_prefetch逻辑用密不透风的预取最大限度让计算掩盖通信。Parallelize 逻辑——注意专家 FSDP 切分时对shard_placement_fn的巧妙处理当efsdp size × ep_degree超过专家总数时专家权重无法沿 dim-0 继续切分此时自动退化为沿 dim-1hidden切分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_forwardreshard_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_forwardreshard_after_forward, )Prefetch 逻辑——前向从 Embedding 出发逐层预取MoE 层预取[next_block, next_block.moe.experts]最后一层预取[norm, output]反向则完全对称地从 Output 逆推回 Embeddingtransformer_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 CPTP/EP/PP 仍在其未来计划中FSDP 通过AutoModelForCausalLM.from_pretrained()自动读取 HuggingFace 架构信息无需权重格式转换这降低了二次开发的适配成本见 slime FSDP 后端。权重同步闭环EP 专家梯度各自优化后最终仍要通过update_weights_from_tensorhandle 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 展示了 EPFSDP2 的干净切分流程与手写 prefetchAutomodel 展示了 DeepEP 的handle状态管理与 autograd 封装TorchTitan 则将前向/反向全链路 prefetch 做到极致——三者共同构成了一幅在 FSDP 上二次开发 EP的完整工程地图。赞分享文档教程人工智能大模型RLHF【免费下载链接】Awesome-ML-SYS-TutorialMy learning notes for ML SYS.项目地址https://gitcode.com/gh_mirrors/aw/Awesome-ML-SYS-Tutorial点击查看免费下载上一篇构建多语言表单jQuery Validation国际化messages文件完全指南下一篇大模型面试笔记从开源协作到知识共享的思考创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表