Prhub

#7085 [veomni] feat: EP-aware sharded delta export (fused expert stacks)

原始 PR 作者 ChangyiYang 合并时间 2026-07-29 11:10 文件变更 9 提交数 9 评论 18 代码增减 +477 / -49

执行摘要

VeOmni EP 感知分片 delta 导出,大幅加速 MoE 权重同步

VeOmni后端使用融合专家堆栈(如gate_up_proj)并结合EP+FSDP2混合并行,原有的分片delta引擎无法正确处理其shard geometry,导致无法使用delta_sharded模式进行高效权重同步。本PR填补了这一空白,使得VeOmni用户也能享受分片delta带来的通信节省。PR body明确指出这是review-requested split的第二部分(引擎核心+FSDP已移至#7144),在本PR中“wires veomni FSDP2+EP into the sharded delta sync end to end”。

值得仔细阅读。特别关注:(1) 如何通过handler自身探测slot表(enumerate_hf_slots)来避免手动维护slot表;(2) NaN sentinel的设计使得每个rank可以独立完成delta转换,无需rank 0全量重建;(3) BlockPlacement如何统一表达EP手动分割和FSDP DTensor分割。

讨论亮点
  • wuxibin89要求将引擎核心重构与VeOmni集成分离为两个PR(#7144和本PR),作者照做。
  • gemini-code-assist[bot]指出gather_dense_blocks_to_rank0dist.gather的dst参数需使用group-relative rank 0,而非dist.get_global_rank
  • wuxibin89担心全量materialize专家内存消耗(如kimi-k2.5的gate_up_proj 21GB),作者澄清该实现使用分片方式避免materialize。
  • wuxibin89询问snapshot是否仅在第一部seed,作者解释hf_delta_export在每次sync时都会刷新snapshot(diff后snap.copy_(local)),无需额外seed。
  • wuxibin89建议将get_per_tensor_param_shard实现移至utils.py以保持简洁,作者在后续commit中重构。

实现拆解

  1. 重构MoE参数转换器:将_map_moe_params_commondefault_moe_param_handlerep_rank参数改为expert_id_base,使其支持任意起点索引,从而兼容EP局部快照和全局全量。
  2. 新增VeOmni专用converter machinery:在verl/workers/engine/veomni/utils.py中实现slot表自动枚举(enumerate_hf_slots)、单行转换(convert_row_to_hf)、delta entry构建器(hf_entry_converter)以及顶层导出生成器(veomni_shard_export),支持融合专家堆栈到HF独立命名的映射。
  3. 扩展spec层:在ShardSpec中新增contributes字段,用于处理HSDP副本维度;将derive_placement重命名为derive_dtensor_placement并明确其职责;同时修改hf_delta_export,使其在spec.place已显式设置时直接使用exporter提供的placement,而不经过DTensor推导。
  4. 在VeOmni引擎中接入:替换transformer_impl.pyget_per_tensor_param_shardNotImplementedError为实际实现,通过调用veomni_shard_export获取带Spec的shard,并实现_hf_delta_entry方法,对含converter的spec分派到hf_entry_converter,否则回退到父类的DTensor逻辑。
  5. 配套测试与文档:在test_sharded_delta.py中新增test_hf_delta_export_converter_param等测试,覆盖converter参数的delta导出、非平凡block placement、以及HSDP replica rank锁步行为;更新delta_weight_sync.md文档。
文件 模块 状态 重要度
verl/workers/engine/veomni/utils.py VeOmni 引擎 modified 9.21
verl/workers/engine/veomni/transformer_impl.py VeOmni 引擎 modified 7.88
tests/checkpoint_engine/test_sharded_delta.py 分片 delta modified 7.82
verl/workers/engine/spec.py 规格层 modified 7.58
verl/workers/engine/utils.py 引擎工具 modified 6.38
docs/advance/delta_weight_sync.md 文档 modified 2.54

关键符号

_map_moe_params_common default_moe_param_handler enumerate_hf_slots convert_row_to_hf hf_entry_converter veomni_shard_export _is_ep_param get_per_tensor_param_shard _hf_delta_entry derive_dtensor_placement hf_delta_export

关键源码片段

verl/workers/engine/veomni/transformer_impl.py dependency-wiring

引擎入口集成,将 get_per_tensor_param_shard 从 NotImplementedError 改为实际实现,新增 _hf_delta_entry 进行 converter dispatch。

def get_per_tensor_param_shard(self, **kwargs):
    """Yield each rank's *local* shard with its ShardSpec."""
    from .utils import veomni_shard_export
​
    # 判断是否需要手动 offload(CPUOffloadPolicy 下不需要)
    manual_offload = not getattr(self, "_uses_fsdp2_cpu_offload_policy", False)
    if manual_offload:
        load_veomni_model_to_gpu(self.module)
    # 调用专用导出函数获取生成器和元数据
    gen, meta = veomni_shard_export(self.module)
​
    def _with_offload_back():
        """包装生成器,在迭代结束后执行 offload(如果适用)"""
        yield from gen
        if manual_offload and self._is_offload_param:
            offload_veomni_model_to_cpu(self.module)
​
    return _with_offload_back(), meta
​
​
def _hf_delta_entry(self, name, spec, place, lidx, lval):
    """veomni 的每参数 delta entry 构建:对含 converter 的 fused expert 参数使用
    hf_entry_converter,否则回退到父类的 DTensor 处理。"""
    from ..spec import BlockPlacement
    from .utils import NO_SLOTS_MSG, hf_entry_converter
​
    # 如果 spec 携带了完整的 converter 信息(to_hf_chunk + hf_slots),
    # 并且 place 是 BlockPlacement,则使用 hf_entry_converter 构建 entry
    if spec.to_hf_chunk is not None and isinstance(place, BlockPlacement) and spec.hf_slots is not None:
        return hf_entry_converter(name, spec, place, lidx, lval)
    # 如果只有 to_hf_chunk 但没有 hf_slots,则无法在发送端完成转换,报错
    if spec.to_hf_chunk is not None:
        raise NotImplementedError(f"{name}: {NO_SLOTS_MSG}")
    # 其余参数走父类的 DTensor identity 处理
    return super()._hf_delta_entry(name, spec, place, lidx, lval)

评论区精华

PR 拆分请求 设计

wuxibin89 要求将 DeltaShardedCheckpointEngine 重构和 FSDP 支持与 VeOmni EP 导出分离为两个 PR。

结论:作者将引擎核心 +FSDP 部分移至 #7144,本 PR 仅包含 VeOmni EP 导出(叠加在 #7144 之上)。 · 已解决

dist.gather 目标 rank 应为 group-relative 而非 global 正确性

gemini-code-assist[bot] 指出在 gather_dense_blocks_to_rank0 中使用 dist.get_global_rank(group, 0) 会导致 subgroup 下运行时错误,应直接使用 0。

结论:接受建议,修改为 group-relative rank 0。 · 已解决

全量 materialize 专家内存风险 性能

wuxibin89 担心 gather_dense_blocks_to_rank0 会在单 rank 上 materialize 全量专家,造成内存爆炸(如 kimi-k2.5 的 gate_up_proj 21GB)。

结论:作者澄清该实现使用分片方式,每个 rank 仅持有自己的 block,不会 materialize 全量。 · 已解决

Snapshot 刷新时机 正确性

wuxibin89 询问 snapshot 是否仅在第一部 seed,后续 sync 不会刷新 snapshot。

结论:作者解释 hf_delta_export 在每次 sync 时都会刷新 snapshot(diff 后 snap.copy_(local)),所以每个 sync 后 snapshot 都对应最新状态,无需额外 seed。 · 已解决

get_per_tensor_param_shard 实现位置 style

wuxibin89 建议将 get_per_tensor_param_shard 的实现移到 veomni/utils.py 以保持 transformer_impl.py 简洁。

结论:作者采纳建议,在 50b488bd 中将 DTensor+EP 声明逻辑移至 utils.py 的 veomni_shard_export。 · 已解决

风险与影响

依赖风险:本PR基于#7144的引擎核心重构,若#7144有回归会影响本PR。
正确性风险:converter machinery涉及EP+FSDP2混合并行placement推导,若BlockPlacement计算错误可能导致权重损坏或训练发散。
性能风险:每次sync中增加的converter调用(每参数NaN probe和slot lookup)可能引入额外开销,但从benchmark看收益远超成本。
兼容性风险:hf_delta_export的dispatch修改可能影响其他后端(如纯FSDP),需确保回归测试覆盖。
内存风险:虽然避免全量materialize,但NaN probe需要构造临时tensor,对于超大MoE仍应注意上限。

用户影响:VeOmni后端用户现在可使用delta_sharded模式,获得4-21倍的权重同步加速。
系统影响:分片delta引擎的能力从仅支持FSDP扩展到支持EP+FSDP混合并行,为后续支持其他后端(如mcore)奠定基础。
团队影响:完成了重构路线图中的重要一步,使VeOmni能充分利用delta增益。

依赖上游重构 (#7144) EP+FSDP 混合并行 placement 复杂 新增 converter 可能引入回归 NaN Probe 临时内存开销

关联 Issue

未识别关联 Issue

当前没有检测到明确关联的 Issue 链接,后续同步到相关引用后会出现在这里。

完整报告

参与讨论