feat: qwen35 share expert and deepep combine overlap - #1343
Open
HongminTan wants to merge 2 commits into
Open
Conversation
LLLLKKKK
previously requested changes
Aug 28, 2026
LLLLKKKK
left a comment
Collaborator
There was a problem hiding this comment.
AI Code Review - PR #1343
Status: BLOCKING
Summary: P0/0 · P1/2 · P2/7 · P3/7
Reviewed: commit 40f6439c08a9 · 2026-08-28 12:29 UTC+8
Blocking Issues
P1
- 默认全量生效的 EP 融合归约与副流重叠缺少真实多卡数值覆盖 @
rtp_llm/models_py/modules/factory/fused_moe/impl/cuda/routers/deepep_low_latency_router.py:104- 建议:按
impl/rocm/test/test_generic_moe_allreduce.py的既有形态补一条 2 rank 用例(routers/test/BUILD已有GPU_COUNT: "2"脚手架):对同一输入分别跑 legacy(_finalize_post_tp_gather+all_reduce(shared),两次集合)与新路径(_finalize_row_scatter+ 单次all_reduce),断言集合次数 2→1 且输出在容差内一致。若受限于多卡资源,请补一条同时满足 DeepEP low-latency + shared expert +ffn_tp_size>1+ep_size>1的 smoke case,并在 PR description 中点明 case 名称,使默认切换具备可追溯的数值防线。
- 建议:按
finalize()新增的 row-scatter 分派分支与能力声明在 CI 中零覆盖 @rtp_llm/models_py/modules/factory/fused_moe/impl/cuda/routers/deepep_low_latency_router.py:364- 建议:补一条
self.assertTrue(_make_router().supports_row_scatter_finalize)对齐同目录先例,防止该能力被误删后use_ep_unified_allreduce静默退回旧双集合路径而测试全绿;并用现有 stub self + patch_normal_finalize覆盖finalize()两条分支(纯 CPU 即可,参考pure_tp_router_skip_allreduce_test.py:88直接调finalize的写法):断言 (a) 给row_scatter_target时返回对象就是传入缓冲区且未发生 all_gather,(b) 不给时走_finalize_post_tp_gather,(c) 两条路径都把self._handle复位为 None。
- 建议:补一条
Non-blocking Suggestions
P2
- 归约合并与跨流重叠没有运行时回滚开关,且无法分别关闭 @
rtp_llm/models_py/model_desc/generic_moe.py:165- 建议:补一个 env 或
MoeConfig开关:使supports_row_scatter_finalize返回 False 即可整体退回已验证的use_ep_shared_allreduce路径;再拆一个更细的开关(例如强制shared_expert_stream=None),使跨流重叠可单独关闭而保留一次归约。默认可继续开启,并在 PR description 中记录默认值与回滚方式。
- 建议:补一个 env 或
- 层未校验 router 返回的即为自己提供的缓冲区,共享专家分支可能被静默丢弃 @
rtp_llm/models_py/model_desc/generic_moe.py:336- 建议:在
all_reduce之前补一行零成本自校验,例如assert experts_output is row_scatter_target, "row-scatter router must return the caller's buffer";这样任何未来 router 的实现偏差都会在首次前向立刻暴露,而不是退化为静默的数值错误。
- 建议:在
supports_row_scatter_finalize契约未写明必须 joinrow_scatter_ready,且缺少统一 join helper @rtp_llm/models_py/modules/factory/fused_moe/defs/fused_moe.py:135- 建议:在
defs/fused_moe.py增加共享 helper(例如join_row_scatter_producer(extra_finalize_args, target))内部完成torch.cuda.current_stream(target.device).wait_event(...),并要求所有 row-scatter 实现在写入缓冲区前第一步调用它;同时把「必须先 join 再触碰缓冲区,即使本 rank 切片为空」与现有两条要求并列写进 docstring,并说明为何 wait 放在 combine 之后而非FusedMoe侧(那样会丧失重叠收益,这个取舍值得留档)。若希望协议不可误用,可把 target 与 event 合并成一个 NamedTuple 传递,使「拿到 target 就必然看到 event」在类型层面成立。
- 建议:在
- capture 用例中「多图复用同一 stream/event」的断言恒为真,给出虚假覆盖信号 @
rtp_llm/models_py/model_desc/test/generic_moe_allreduce_test.py:434- 建议:改为对被测对象的实际产出取证:在循环内把
_ensure_shared_expert_stream(self.device)的返回值(而非 stub 属性)加入forked_on并断言len(_SHARED_EXPERT_STREAMS) == 1,或让_fork_and_join返回ready后收集其id,或对torch.cuda.Stream构造计数。若该不变量已由test_one_stream_per_device_is_reused充分覆盖,则直接删除这条不可失败的断言并相应改名,只保留能真正失败的多图 capture + replay 数值校验。
- 建议:改为对被测对象的实际产出取证:在循环内把
- 副流的真实装配与端到端 fork→join 接线在所有用例中零覆盖 @
rtp_llm/models_py/model_desc/test/generic_moe_allreduce_test.py:366- 建议:给
_make_layer增加device参数以支持 CUDA 权重,并补一条 GPU 用例:断言layer.shared_expert_stream is _ensure_shared_expert_stream(self.device)、layer.shared_expert_ready是真实torch.cuda.Event、且fused_moe.call_args.kwargs["row_scatter_ready"] is layer.shared_expert_ready;再以真实DenseMLP作为shared_expert,与「强制shared_expert_stream=None的串行结果」逐元素比对,覆盖 eager 与 capture/replay 两种模式。
- 建议:给
- 副流不可用时重叠静默降级,且实际归约策略仅在 DEBUG 日志部分可见 @
rtp_llm/models_py/model_desc/generic_moe.py:189- 建议:把
shared_expert_stream is not None与use_ep_shared_allreduce并入现有日志字段;当use_ep_unified_allreduce为真却拿不到副流时打一条logger.warning显式说明原因(权重不在 CUDA 设备上);或在模型构建完成后以 INFO 统计一次启用重叠的 MoE 层数,让配置生效情况在默认日志级别下可验证,而不是只能靠 profiling 反推。
- 建议:把
- 测试替身用硬编码与无 spec mock,削弱了三个生产边界 @
rtp_llm/models_py/model_desc/test/generic_moe_allreduce_test.py:69- 建议:补一条不 patch
MoEConfigAdapter的窄测试,用真实parallelism_config(令get_attn_tp_size()与get_ffn_tp_size()返回不同值)断言FusedMoeDataRouter.tp_collective_size取的是 attn 视图;或在 fixture 中由get_attn_tp_size()派生tp_collective_size,使该前提在测试内显式成立而非被硬编码。对集合通信与 DenseMLP 两个边界改用autospec=True/create_autospec(DenseMLP, instance=True),让签名不匹配在单测阶段即报错。
- 建议:补一条不 patch
P3
- row-scatter 的形状/dtype 契约仅由 assert 保障,且负向覆盖不完整 @
rtp_llm/models_py/modules/factory/fused_moe/impl/cuda/routers/deepep_low_latency_router.py:314- 建议:把这三条契约从
assert改为显式raise ValueError(保留当前 f-string 诊断信息),或在FusedMoe.forward的入参校验块中前移row_scatter_target的 shape/dtype/device 检查,与已有两条ValueError保持同一层级;同时把该负向用例扩展为subTest矩阵,分别构造 hidden 维不一致与 dtype 不一致的输入并断言各自文案。
- 建议:把这三条契约从
- finalize 中 assert 失败会残留 router handle,后续请求报无关错误 @
rtp_llm/models_py/modules/factory/fused_moe/impl/cuda/routers/deepep_low_latency_router.py:319- 建议:把
self._handle = None放进try/finally(或在finalize入口把句柄取到局部变量后立即清空实例字段),使无论走哪条分支、是否抛异常,router 都回到可复用状态,assert 的失败信息也能如实反映首次故障点;若认为该复位超出本 PR 范围,请在generic_moe.py的except注释中明确记录「router 内部状态未回滚」这一已知限制。
- 建议:把
- 契约测试用
object()冒充torch.cuda.Event,与类型声明不符 @rtp_llm/models_py/modules/factory/fused_moe/defs/test/fused_moe_allreduce_contract_test.py:156- 建议:
torch.cuda.Event()的底层 event 是惰性创建的,无 CUDA 设备也可构造,建议直接改用真实 Event 以对齐类型声明;或至少改成Mock(spec=torch.cuda.Event),使任何非 Event 接口调用立刻失败,同时assertIs身份断言写法完全不变。
- 建议:
- row-scatter 契约测试缺两个边界组合,且两个契约合并在同一方法内 @
rtp_llm/models_py/modules/factory/fused_moe/defs/test/fused_moe_allreduce_contract_test.py:142- 建议:拆成两个测试方法或用
subTest参数化(target, ready, expect_keys),补齐「target 有、ready 无」组合(断言含ROW_SCATTER_TARGET_ARG且不含ROW_SCATTER_READY_ARG);若保留单方法写法,至少在第二次调用前router.finalize.reset_mock()。并为互斥组合补一条用例,或先在FusedMoe.forward中对该组合抛ValueError再补测。
- 建议:拆成两个测试方法或用
finalize仅校验位置参数个数,未固定参数顺序 @rtp_llm/models_py/modules/factory/fused_moe/defs/test/fused_moe_allreduce_contract_test.py:20- 建议:改用
create_autospec(FusedMoeDataRouter, instance=True)让位置/关键字绑定由 spec 校验,并在提取时按参数名读取(绑定后取extra_finalize_args),替代按下标手工提取;这样参数顺序漂移在单测阶段即暴露,也无需在测试内重复维护「五个位置参数」这一弱约定。
- 建议:改用
_SHARED_EXPERT_STREAMS缺少清理入口,测试直接清空生产全局 @rtp_llm/models_py/model_desc/generic_moe.py:39- 建议:提供显式清理函数(例如
clear_shared_expert_streams())供拆卸路径与测试调用;测试侧改用patch.dict(..., {}, clear=True)或addCleanup还原快照,避免全局状态跨用例泄漏。若确认 layer 构造只发生在单线程加载阶段,请在注释中写明该前提,说明为何不需要加锁。
- 建议:提供显式清理函数(例如
- 三个归约策略布尔量互相纠缠且依赖赋值顺序 @
rtp_llm/models_py/model_desc/generic_moe.py:172- 建议:至少把顺序依赖显式化(例如先算出
ep_shared_eligible再由它派生两个 flag),使赋值顺序不再隐含正确性;若后续还要新增策略,可考虑把归约策略收敛为一个枚举(如EP_UNIFIED / TP_UNIFIED / EP_SHARED_ONLY / LOCAL_ONLY)在__init__一次性计算、forward单点分派,注释里的策略表可直接对应枚举值。
- 建议:至少把顺序依赖显式化(例如先算出
Checklist Findings (14 fail / 44 total)
General Principles Checklist
- [6.1] Architecture — 分层边界:新概念在正确层级,不泄漏内部 → issue ``_SHARED_EXPERT_STREAMS
缺少清理入口,测试直接清空生产全局
_`_SHARED_EXPERT_STREAMS` 是模块级可变全局,`_ensure_shared_expert_stream` 的 get/create/set 之间无锁,也没有任何清理入口:拆卸路径不会清空,模型重建或多模型共存时会静默复用旧 stream,stream 在进程生命周期内不释放。测试只能自己在 `generic_moe_allreduce_test.py:363` 的 `setUp` 里直接 `SHARED_EXPERT_STREAMS.clear()` 且无 `tearDown` 还原——这既耦合了私有全局,又与该全局注释声明的「每设备一条、必须在首次 captured forward 前存在」不变量相悖,为同 target 内其它用例引入隐式顺序依赖(当前其它用例均用 CPU 权重、无该全局读者,实际影响有限)。 - [6.1] Architecture — 可观测性:日志/指标/超时可操作、非噪声 → issue
副流不可用时重叠静默降级,且实际归约策略仅在 DEBUG 日志部分可见
_ensure_shared_expert_stream对非 CUDA device 直接返回 None(第 45 行),此时shared_expert_ready也不创建,_shared_expert_partial退回同流串行——本 PR 的全部重叠收益消失,却不产生任何日志或告警。若某条权重加载路径在__init__时权重仍在 CPU/meta 上(随后才.to(cuda)),use_ep_unified_allreduce仍为 True 而重叠被永久静默关闭。紧随的logger.debug(第 193-203 行)只打印use_unified_tp_allreduce与use_ep_unified_allreduce,既不含use_ep_shared_allreduce,也不含副流是否真正取到,且整段被isEnabledFor(logging.DEBUG)包住,线上 INFO 级别无法区分「重叠生效」与「重叠静默失效」。 - [6.1] Architecture — 回滚路径:风险行为存在运维回滚手段 → issue
归约合并与跨流重叠没有运行时回滚开关,且无法分别关闭
use_ep_unified_allreduce的五个合取项全部来自并行度与 router 能力,全仓搜索确认无任何 env 或MoeConfig字段可关闭。一旦线上出现精度漂移、流并发或 SM 争抢导致的 TPOT 回退,唯一手段是回滚镜像。更关键的是「合并 all-reduce」(纯代数等价、风险低)与「跨流重叠」(引入并发、风险明显更高)被绑定在同一个谓词上:出问题时无法只关掉后者而保留前者的收益。 - [6.1] Architecture — 状态不变量:创建/更新/失败/重试/回滚路径有效 → issue ``_SHARED_EXPERT_STREAMS
缺少清理入口,测试直接清空生产全局
_`_SHARED_EXPERT_STREAMS` 是模块级可变全局,`_ensure_shared_expert_stream` 的 get/create/set 之间无锁,也没有任何清理入口:拆卸路径不会清空,模型重建或多模型共存时会静默复用旧 stream,stream 在进程生命周期内不释放。测试只能自己在 `generic_moe_allreduce_test.py:363` 的 `setUp` 里直接 `SHARED_EXPERT_STREAMS.clear()` 且无 `tearDown` 还原——这既耦合了私有全局,又与该全局注释声明的「每设备一条、必须在首次 captured forward 前存在」不变量相悖,为同 target 内其它用例引入隐式顺序依赖(当前其它用例均用 CPU 权重、无该全局读者,实际影响有限)。 - [6.1] Architecture — 错误语义:fail-fast/retry/fallback/silent 行为显式 → issue
finalize 中 assert 失败会残留 router handle,后续请求报无关错误
_finalize_row_scatter新增的三处 assert(dim / shape+dtype /combined_x.size(0) == slice_size)均在self._handle = None(第 376 行)之前执行。任一触发后_handle仍持有上一次 dispatch 的句柄,下一次请求会撞上prepare里的assert self._handle is None(第 223 行),报出与真实原因完全无关的错误,router 实例事实上永久不可用——一次可恢复的单请求失败被放大成引擎级故障。该闩锁是既有行为(旧_finalize_post_tp_gather也有 assert 在复位之前),但本 PR 扩大了失败面;generic_moe.py:327新增的except只补做 stream join、未复位 router 状态,注释却容易让读者误以为失败路径已完整。 - [6.1] Software Engineering — LSP:子类/重写保持基类契约 → issue ``supports_row_scatter_finalize
契约未写明必须 joinrow_scatter_ready`,且缺少统一 join helper`
该 docstring(第 136-146 行)详细要求「累加而非覆盖」「只写本 rank 拥有的行」,却完全没提 `row_scatter_ready`。副流 fork 由 `generic_moe.py:315` 打开,唯一成功路径的 join 却内联在另一模块的 `deepep_low_latency_router.py:311-312`;`fused_moe.py:344-347` 只在事件非空时透传,从不校验是否被消费,router 侧又用 `.get()` 宽松读取。第二个实现者若忽略该键,eager 模式下是静默跨流数据竞争(读到未完成的 shared partial),只有捕获模式才以 fork 未闭合暴露。对照同类能力 `skip_tp_allreduce` 在 `fused_moe.py:80` 提供了共享判定函数 `should_skip_tp_allreduce`,本次 join 逻辑却无共享入口。 - [6.1] Software Engineering — OCP:本地扩展点优先于修改中心逻辑 → issue
三个归约策略布尔量互相纠缠且依赖赋值顺序
use_ep_shared_allreduce的定义里包含not self.use_ep_unified_allreduce,因此必须在后者之后赋值,这个顺序约束只靠代码位置隐式保证;forward中四条分支(ep_unified / unified_tp / ep_shared / fallback)的互斥性也只能靠人工阅读第 140-164 行的 25 行 ASCII 注释块确认。再加入第四种策略时,这张布尔网的组合正确性会难以维持。 - [6.1] Software Engineering — SRP:模块/类职责单一 → issue
三个归约策略布尔量互相纠缠且依赖赋值顺序
use_ep_shared_allreduce的定义里包含not self.use_ep_unified_allreduce,因此必须在后者之后赋值,这个顺序约束只靠代码位置隐式保证;forward中四条分支(ep_unified / unified_tp / ep_shared / fallback)的互斥性也只能靠人工阅读第 140-164 行的 25 行 ASCII 注释块确认。再加入第四种策略时,这张布尔网的组合正确性会难以维持。 - [6.1] Tests — 分布式/跨平台变更有对应覆盖 → issue
副流的真实装配与端到端 fork→join 接线在所有用例中零覆盖
_make_layer的W.moe_w1/moe_w2恒为 CPU 张量(第 59-60 行),故generic_moe.py:190的_ensure_shared_expert_stream(self.w1.device)必返回 None、shared_expert_ready亦为 None,即使在 H20 上亦如此——第 306 行正是断言这一点。唯一触达副流的SharedExpertSideStreamCaptureTest._stub_layer用SimpleNamespace冒充 layer、shared_expert是lambda x, skip_allreduce: x * 3.0,join 也由测试体自行调用而非经 router。于是「构造期按w1.device建流建 event → forward 透传 → router 在_finalize_row_scatterjoin」这条生产链路,以及真实量化DenseMLPGEMM 与 routed 链路并发的数值等价性,均无任何用例执行 - [6.1] Tests — 新逻辑有聚焦单测 + 相关集成/smoke 测试 → issue ``finalize
仅校验位置参数个数,未固定参数顺序
_`extract_extra_finalize_args` 通过 `len(args) != 5` 手工校验个数后取 `args[4]`。生产调用为 `router.finalize(combine_payload, expert_topk_weights, expert_topk_ids, apply_router_weight_on_input, finalize_args)`(`fused_moe.py:349-355`)。由于 `router` 是 `Mock(spec=FusedMoeDataRouter)`(spec 只校验属性存在性、不校验调用签名),若 `forward` 误把 `topk_weights` 与 `topk_ids`、或 `apply_router_weight_on_input` 的位置互换,个数仍为 5,本文件全部断言依旧通过,而所有真实 router 都会拿到错位参数。 - [6.1] Tests — 边界 case 覆盖(空、单元素、最大值) → issue
row-scatter 契约测试缺两个边界组合,且两个契约合并在同一方法内
FusedMoeRowScatterTargetTest只覆盖「两者都不传」与「target 与 ready 都传」。但生产上GenericMoeLayer在没有副流时恒产生「只传 target、不传 ready」(generic_moe_allreduce_test.py:306正是这一形态),router 会走.get(ROW_SCATTER_READY_ARG)为 None 的免 join 分支,契约层未固化。另外skip_tp_allreduce=True与row_scatter_target属语义互斥组合,forward既无校验也无用例,仅靠基类两个能力属性默认 False 而暂不可达。同时第 142-166 行在一个方法里先后调用两次 forward、依赖call_args只反映最后一次来断言,两个方向共享同一失败信号,与同文件第 72 行已用subTest的风格也不一致。
RTP-LLM Checklist
- [I] 代码质量 — 同一功能用统一工具函数 → issue ``supports_row_scatter_finalize
契约未写明必须 joinrow_scatter_ready`,且缺少统一 join helper`
该 docstring(第 136-146 行)详细要求「累加而非覆盖」「只写本 rank 拥有的行」,却完全没提 `row_scatter_ready`。副流 fork 由 `generic_moe.py:315` 打开,唯一成功路径的 join 却内联在另一模块的 `deepep_low_latency_router.py:311-312`;`fused_moe.py:344-347` 只在事件非空时透传,从不校验是否被消费,router 侧又用 `.get()` 宽松读取。第二个实现者若忽略该键,eager 模式下是静默跨流数据竞争(读到未完成的 shared partial),只有捕获模式才以 fork 未闭合暴露。对照同类能力 `skip_tp_allreduce` 在 `fused_moe.py:80` 提供了共享判定函数 `should_skip_tp_allreduce`,本次 join 逻辑却无共享入口。
Python Static-First Checklist
- [P.G] 测试规范 — mock/fake/stub 不得替代本次声称覆盖的生产边界 → issue ``finalize
仅校验位置参数个数,未固定参数顺序
_`extract_extra_finalize_args` 通过 `len(args) != 5` 手工校验个数后取 `args[4]`。生产调用为 `router.finalize(combine_payload, expert_topk_weights, expert_topk_ids, apply_router_weight_on_input, finalize_args)`(`fused_moe.py:349-355`)。由于 `router` 是 `Mock(spec=FusedMoeDataRouter)`(spec 只校验属性存在性、不校验调用签名),若 `forward` 误把 `topk_weights` 与 `topk_ids`、或 `apply_router_weight_on_input` 的位置互换,个数仍为 5,本文件全部断言依旧通过,而所有真实 router 都会拿到错位参数。 - [P.G] 测试规范 — 数据驱动测试用 pytest.mark.parametrize → issue
row-scatter 契约测试缺两个边界组合,且两个契约合并在同一方法内
FusedMoeRowScatterTargetTest只覆盖「两者都不传」与「target 与 ready 都传」。但生产上GenericMoeLayer在没有副流时恒产生「只传 target、不传 ready」(generic_moe_allreduce_test.py:306正是这一形态),router 会走.get(ROW_SCATTER_READY_ARG)为 None 的免 join 分支,契约层未固化。另外skip_tp_allreduce=True与row_scatter_target属语义互斥组合,forward既无校验也无用例,仅靠基类两个能力属性默认 False 而暂不可达。同时第 142-166 行在一个方法里先后调用两次 forward、依赖call_args只反映最后一次来断言,两个方向共享同一失败信号,与同文件第 72 行已用subTest的风格也不一致。
Strengths
- 归约代数被严格证明而非口头声明:
deepep_low_latency_row_scatter_test.py:113以all_gather(routed) + all_reduce(shared)为基线对拍融合结果,token 数覆盖 1/3/5/8/9(整除、ragged 尾部、尾部 rank 完全无 token),并用assertEqual(routed_full.shape[0], original_num_tokens)顺带证明切片确实构成划分;target 取非零的gate * shared_partials[rank],能同时抓住「add_退化为覆盖」与「写入他 rank 行」两类缺陷。 _tp_token_slice(deepep_low_latency_router.py:107)把切片算术从_prepare_pre_tp_slice抽为单一实现,使「dispatch 切法」与「scatter 回写行」不可能漂移,并成为该不变量的唯一来源;ragged 尾部因此不再需要旧路径的torch.empty等长补齐(:283-288)。_finalize_row_scatter把wait_event放在所有 assert 与写入之前并注释理由(调用方归约全部行,空切片 rank 同样必须 join);deepep_low_latency_row_scatter_test.py:83在 patch 的wait_event回调里断言 buffer 仍全零,是真正会失败的时序断言而非调用计数,且覆盖了尾部空 rank。- 跨流内存生命周期双向到位:
hidden_states.record_stream(stream)防止上层解引用后 allocator 复用仍被副流读取的块,shared_partial.record_stream(current)防止副流分配块在主流仍在写时被回收,注释说明了各自原因。 SharedExpertSideStreamCaptureTest选择真实torch.cuda.CUDAGraph而非 mock 验证 fork 闭合(未闭合只在 capture end 报错),并覆盖「多 batch size 共享同一 graph pool」这一解码实际形态,replay 前hidden.fill_改值以区分真实重算与残留缓冲;_CUDA_GRAPH_READY显式排除 HIP 并注明理由,使generic_moe_allreduce_test_rocm目标不会误跑 CUDA-only 分支。FusedMoe.forward对能力不匹配与参数不自洽在router.prepare之前 fail-fast,三条ValueError用例均额外断言prepare/execute未被调用;forward又在 routed 链路抛异常时补做wait_event,generic_moe_allreduce_test.py:309覆盖了这条最易漏掉的失败路径。test_layer_only_passes_keywords_fused_moe_accepts(generic_moe_allreduce_test.py:329)用inspect.signature(FusedMoe.forward)反查 kwargs,精准堵住「MagicMock 接受任意关键字导致 layer 多传参数直到线上才炸」的盲区,四条归约分支全覆盖。- 归约策略决策表以结构化注释集中写在
generic_moe.py:140-164,显式区分「partial 可合并」与「已完整不可合并」,把本次改动最易出错的前提(已完整分支参与合并归约会被乘 tp_size)留档,后续维护者不必反推。
LLLLKKKK
reviewed
Aug 28, 2026
HongminTan
force-pushed
the
qwen35-share-expert-overlap
branch
from
August 31, 2026 07:25
40f6439 to
afdeb08
Compare
LLLLKKKK
reviewed
Aug 31, 2026
LLLLKKKK
left a comment
Collaborator
There was a problem hiding this comment.
AI Code Review - PR #1343
Status: LGTM
Summary: P0/0 · P1/0 · P2/8 · P3/13
Reviewed: commit afdeb0846375 · 2026-08-31 16:08 UTC+8
lgtm ready to ci
Non-blocking Suggestions
P2
- extra_finalize_args 可绕过 row_scatter 的两道能力门禁 @
rtp_llm/models_py/modules/factory/fused_moe/defs/fused_moe.py:344- 建议:按
skip_tp_allreduce对齐:在finalize_args.update()中无条件写入这两个键(为 None 时显式pop),使函数参数成为唯一真源;并在FusedMoeRowScatterTargetTest补两条与test_unsupported_router_cannot_be_bypassed_by_finalize_key对称的负向用例,断言注入的row_scatter_target/row_scatter_ready不会泄漏到router.finalize。
- 建议:按
- supports_row_scatter_finalize 契约未写明必须 join row_scatter_ready,也未写调用方 buffer 义务 @
rtp_llm/models_py/modules/factory/fused_moe/defs/fused_moe.py:135- 建议:把「触碰 buffer 前必须
wait_event(row_scatter_ready),即使本 rank 切片为空」写成与「累加不覆写」同级的硬性要求,并给FinalizeArgs.row_scatter_ready加注释;更稳的做法是把 join 上移到FusedMoe.forward调用router.finalize之前(此时 dispatch 与 expert 计算均已入队,重叠效果不变),让契约由框架强制。同时把调用方义务补进 docstring,并在FusedMoe.forward现有校验块加上is_contiguous()与与hidden_states的存储别名 fail-fast 检查。
- 建议:把「触碰 buffer 前必须
- 层未校验 router 返回的即为自己提供的缓冲区,共享专家分支可能被静默丢弃 @
rtp_llm/models_py/model_desc/generic_moe.py:336- 建议:在
FusedMoe.forward中当row_scatter_target is not None时增加assert output is row_scatter_target,把契约违约从静默错误变为 fail-fast;并把_make_fused_moe的 row-scatter 分支改为合规假实现(对传入 target 执行add_后返回 target),补一条assertIs(output, target),使契约测试不再编码违约行为。
- 建议:在
- 归约合并与跨流重叠没有运行时回滚开关,且无法分别关闭 @
rtp_llm/models_py/model_desc/generic_moe.py:165- 建议:参照仓内已有风格增加默认开启、可关闭的开关(env 或
MoeConfig字段)参与谓词,关闭时自然回落use_ep_shared_allreduce;建议把「row-scatter 融合归约」与「共享专家侧流 overlap」拆成两个独立开关,因为二者失效模式(数值/形状 vs 并发/SM 争用)完全不同。
- 建议:参照仓内已有风格增加默认开启、可关闭的开关(env 或
- 真实 DeepEP combine 未覆盖非整除切片与空切片 @
rtp_llm/models_py/modules/factory/fused_moe/impl/cuda/routers/test/deepep_low_latency_router_test.py:371- 建议:在
_run_deepep_low_latency_router_test中为_check_row_scatter_matches_gather增加至少两组 token 数:一个不被tp_size整除(如num_token_per_rank - 1),一个小于tp_size使尾部 rank 切片为空,让硬断言与slice_size == 0短路在真实 combine kernel 下被验证;若 DeepEP 对 0 行 dispatch 另有约束,也应在该测试中显式暴露。
- 建议:在
- capture 用例中「多图复用同一 stream/event」的断言恒为真 @
rtp_llm/models_py/model_desc/test/generic_moe_allreduce_test.py:434- 建议:让断言观察真实来源:把
_stub_layer()移入循环,使每轮真正经过_ensure_shared_expert_stream,再断言两轮取到的是同一 stream 对象;或在循环内断言_ensure_shared_expert_stream(self.device) is stub.shared_expert_stream。若确认 stub 模式下无法有效验证,宁可删除该恒真断言与其失败消息,避免留下虚假覆盖信号。
- 建议:让断言观察真实来源:把
- 侧流的真实装配与端到端 fork→join 接线在所有用例中零覆盖 @
rtp_llm/models_py/model_desc/test/generic_moe_allreduce_test.py:306- 建议:给
_make_layer加device参数把W.moe_w1/W.moe_w2放到 cuda,在_CUDA_GRAPH_READY保护下新增一条用真实GenericMoeLayer的用例:断言layer.shared_expert_stream is _ensure_shared_expert_stream(device)、layer.shared_expert_ready is not None,并在一次 forward 后断言fused_moe.call_args.kwargs["row_scatter_ready"] is layer.shared_expert_ready,闭合__init__ → forward → router的 event 传递链;保留现有 CPU 用例继续覆盖 None 分支。该 target 已配置 H20 与 MI308X(gpu_count: 1),无需新增 CI 资源。
- 建议:给
- 侧流不可用时重叠静默降级,且实际归约策略仅在 DEBUG 日志部分可见 @
rtp_llm/models_py/model_desc/generic_moe.py:189- 建议:把
self.shared_expert_stream is not None加入 :194 的日志字段,并把最终生效策略从 DEBUG 提升为启动期一次性 INFO;若use_ep_unified_allreduce为真而侧流拿不到,logger.warning一次并说明降级原因(设备类型 / CUDA 不可用),使「策略已启用但重叠未生效」这一状态可被观测。
- 建议:把
P3
- 四种归约策略靠三个互相纠缠且依赖赋值顺序的布尔标志编码 @
rtp_llm/models_py/model_desc/generic_moe.py:172- 建议:把策略收敛为一个枚举或小策略对象,携带「共享分支是否 skip 自身 all-reduce / routed 分支是否接收 scatter target / 集合通信次数」三项属性;
forward只做一次分派,表格注释退化为枚举成员文档,后续新增策略变为新增枚举成员而非修改中心逻辑。
- 建议:把策略收敛为一个枚举或小策略对象,携带「共享分支是否 skip 自身 all-reduce / routed 分支是否接收 scatter target / 集合通信次数」三项属性;
- _finalize_post_tp_gather 未复用同批引入的切片工具 @
rtp_llm/models_py/modules/factory/fused_moe/impl/cuda/routers/deepep_low_latency_router.py:278- 建议:让
_finalize_post_tp_gather直接调用_tp_token_slice,或抽出共享的_tp_token_shard_size(num_tokens),使均分宽度与落位区间只有一个定义点。
- 建议:让
- _ensure_shared_expert_stream 与 dsv4 同名实现近乎逐行重复 @
rtp_llm/models_py/model_desc/generic_moe.py:39- 建议:把 per-device 侧流缓存抽到共享位置(例如
models_py/utils/下的 side-stream helper),由 dsv4 与 generic_moe 共同复用,并顺带提供测试可用的清理入口,避免测试直接清空生产模块全局态。
- 建议:把 per-device 侧流缓存抽到共享位置(例如
- 捕获期无条件调用 record_stream,可能抬高 graph pool 占用 @
rtp_llm/models_py/model_desc/generic_moe.py:256- 建议:评估在捕获期跳过
record_stream(捕获期生命周期由 graph 私有内存池保证,不依赖 allocator 的stream_uses);若确认必须保留,请在注释中写明理由,并给出捕获前后 graph pool 占用的对比数据。
- 建议:评估在捕获期跳过
- all_reduce 的 inplace 用法在两条统一归约路径间不一致 @
rtp_llm/models_py/model_desc/generic_moe.py:365- 建议:统一两条路径的调用形式(pure-TP 分支也传
inplace=True),省掉 symm_mem 快路径上的一次 [T, H] 分配;若有意保持差异,请在注释中说明原因。
- 建议:统一两条路径的调用形式(pure-TP 分支也传
- 真机测试统计的 all_reduce 次数是自证断言 @
rtp_llm/models_py/modules/factory/fused_moe/impl/cuda/routers/test/deepep_low_latency_router_test.py:208- 建议:把
collective_torch.all_reduce一并以wraps方式 patch 计数,让「Group.TP 集合通信从 2 次降到 1 次」成为可失败的断言;若不便 patch,则把 docstring 与断言收敛为「row scatter 路径不再触发 all_gather」,从比较字典中移除 all_reduce 项。
- 建议:把
- row-scatter 的形状/dtype 契约仅由 assert 保障,且负向覆盖不完整 @
rtp_llm/models_py/modules/factory/fused_moe/impl/cuda/routers/deepep_low_latency_router.py:314- 建议:在
DeepEpLowLatencyRowScatterTest补两个小用例:传入 hidden 维不同的 target、传入 dtype 不同(如 bf16 target 配 fp32 combine 输出)的 target,各自用assertRaises(AssertionError)钉住 fail-fast 行为。
- 建议:在
- finalize 中 assert 失败会残留 router handle,后续请求报无关错误 @
rtp_llm/models_py/modules/factory/fused_moe/impl/cuda/routers/deepep_low_latency_router.py:319- 建议:用
try/finally包住 finalize 的后处理,把self._handle = None放进finally(或在 combine 之后立即复位),保证断言失败也把 router 状态复位,使首次抛出的异常即为真实根因。
- 建议:用
- 契约测试用 object() 冒充 torch.cuda.Event,且两个场景挤在同一方法 @
rtp_llm/models_py/modules/factory/fused_moe/defs/test/fused_moe_allreduce_contract_test.py:156- 建议:拆成独立用例或用
subTest区分场景,并补一条「target 有值、ready 为 None」的用例断言 finalize 字典中不含 ready 键;ready改用Mock(spec=torch.cuda.Event),在保持平台中立的同时保留类型契约。
- 建议:拆成独立用例或用
- finalize 断言仅校验位置参数个数,未固定参数顺序 @
rtp_llm/models_py/modules/factory/fused_moe/defs/test/fused_moe_allreduce_contract_test.py:20- 建议:在 helper 中一并断言前四个位置参数的身份(
args[0].fused_expert_output、args[1] is topk_weights、args[2] is topk_ids、args[3] is False),或改用create_autospec/Mock(spec=FusedMoeDataRouter.finalize)之类可校验绑定的替身。
- 建议:在 helper 中一并断言前四个位置参数的身份(
- setUp 清空生产模块全局缓存且不还原 @
rtp_llm/models_py/model_desc/test/generic_moe_allreduce_test.py:363- 建议:改为作用域受限写法:
self.enterContext(patch.dict(_SHARED_EXPERT_STREAMS, {}, clear=True)),或先快照再self.addCleanup还原;更根本的做法是由生产模块提供正规的清理入口,避免测试对生产全局态留下不可逆改动与顺序依赖。
- 建议:改为作用域受限写法:
- 测试在循环体内定义闭包并捕获循环变量 @
rtp_llm/models_py/modules/factory/fused_moe/impl/cuda/routers/test/deepep_low_latency_row_scatter_test.py:132- 建议:改为默认参数绑定(
def wait_event(seen, target=target, waited=waited):),或把单轮断言逻辑抽成接收target的模块级 helper,让绑定关系显式化。
- 建议:改为默认参数绑定(
- 侧流未设优先级,且 PR 缺少热路径实测数据 @
rtp_llm/models_py/model_desc/generic_moe.py:50- 建议:在 PR 描述中补上目标机型(EP+TP 配置)下 small-batch decode 与 prefill 大 batch 的端到端延迟前后对比;若观测到 routed 链被拖慢,考虑用
torch.cuda.Stream(device=..., priority=...)把侧流降为低优先级,或引入类似 dsv4 的 token 阈值保护。
- 建议:在 PR 描述中补上目标机型(EP+TP 配置)下 small-batch decode 与 prefill 大 batch 的端到端延迟前后对比;若观测到 routed 链被拖慢,考虑用
Checklist Findings (18 fail / 44 total)
General Principles Checklist
- [6.1] Architecture — 分层边界:新概念在正确层级,不泄漏内部 → issue
setUp 清空生产模块全局缓存且不还原
setUp直接_SHARED_EXPERT_STREAMS.clear()(:363)修改被测模块的全局缓存,且无 tearDown / addCleanup 还原。该表是「每设备一条侧流」这一 capture 前提的唯一持有者(generic_moe.py:36-52)。当前文件其余用例都用 CPU 权重、不写入该表,故影响尚未显现;但一旦按本报告建议新增构造真实 CUDA layer 的用例,清表会让_ensure_shared_expert_stream为同一设备再发一条流,破坏本类test_one_stream_per_device_is_reused断言的不变量,并使先前已 capture 的 fork 边指向无人 join 的流。 - [6.1] Architecture — 可观测性:日志/指标/超时可操作、非噪声 → issue
侧流不可用时重叠静默降级,且实际归约策略仅在 DEBUG 日志部分可见
_ensure_shared_expert_stream在device.type != "cuda"或 CUDA 不可用时返回 None(:45-46),此时shared_expert_stream/shared_expert_ready保持 None,_shared_expert_partial退化为串行内联执行(:241-248)——功能仍正确但重叠收益全部丢失。紧随其后的日志(:193-203)只打印两个策略标志与并行度,不含侧流是否真的建立,且整段决策只在 DEBUG 级可见。若self.w1在构造期尚未落到 CUDA(延迟加载 / 权重先在 CPU),重叠会整层失效而日志看起来一切正常,性能回退极难定位。 - [6.1] Architecture — 回滚路径:风险行为存在运维回滚手段 → issue
归约合并与跨流重叠没有运行时回滚开关,且无法分别关闭
use_ep_unified_allreduce完全由静态拓扑条件(:165-171)与DeepEpLowLatencyRouter.supports_row_scatter_finalize硬编码return True(router:104-105)推导;全仓检索确认无任何 env 或moe_config开关参与。凡「DeepEP LL router + ep_size>1 + ffn_tp_size>1 + ffn_tp==attn_tp + 有共享专家」的线上配置升级后会一次性同时切换三项行为:归约拓扑与 bf16 累加顺序改变、引入跨流并发、CUDA graph 内新增 fork/join;同时use_ep_shared_allreduce因 :176 被挤掉,旧路径不再可达。仓内同类实现dsv4/moe/shared_expert.py:27,471-474用 env 选择模式(默认 sequential)并带 token 阈值,本路径两者皆无,出问题只能回滚镜像。 - [6.1] Architecture — 状态不变量:创建/更新/失败/重试/回滚路径有效 → issue
setUp 清空生产模块全局缓存且不还原
setUp直接_SHARED_EXPERT_STREAMS.clear()(:363)修改被测模块的全局缓存,且无 tearDown / addCleanup 还原。该表是「每设备一条侧流」这一 capture 前提的唯一持有者(generic_moe.py:36-52)。当前文件其余用例都用 CPU 权重、不写入该表,故影响尚未显现;但一旦按本报告建议新增构造真实 CUDA layer 的用例,清表会让_ensure_shared_expert_stream为同一设备再发一条流,破坏本类test_one_stream_per_device_is_reused断言的不变量,并使先前已 capture 的 fork 边指向无人 join 的流。 - [6.1] Architecture — 错误语义:fail-fast/retry/fallback/silent 行为显式 → issue
finalize 中 assert 失败会残留 router handle,后续请求报无关错误
self._handle = None(:376)位于两条 finalize 分支之后。若_finalize_row_scatter的四条 assert(:314-322)任一触发(例如 combine 返回行数与切片不符),异常在 :367 向上抛出而self._handle仍保留上次 dispatch 的值;下一次请求进入prepare时会撞上assert self._handle is None(:223),报出与真实故障完全无关的错误信息,掩盖首发根因。gather 路径同样存在此结构,但本次新增的硬断言把触发面扩大了,combined_x.size(0) == slice_size正是最可能在 ragged token 下触发的一条。 - [6.1] Quality — PR description 说明动机与设计 → issue
侧流未设优先级,且 PR 缺少热路径实测数据
torch.cuda.Stream(device=torch.device("cuda", index))用默认优先级创建。改动后共享专家的两次 GEMM 与 routed 链(DeepEP dispatch/combine + grouped GEMM)并发争抢 SM,而 routed 链正是 decode 的关键路径。改动动机是 overlap 提速,但 diff 与新增用例只验证正确性,无任何 decode 端到端延迟前后对比;prefill 大 batch 或 small-batch decode 下共享专家抢占 SM 反而拖慢 routed 链完全可能,仅从代码无法判断净收益方向。对比dsv4/moe/shared_expert.py:471-474带 token 阈值(默认 4096)在大 batch 下退回串行,本路径无对应保护。 - [6.1] Software Engineering — DRY:重复非平凡逻辑被抽取或显式复用 → issue
_ensure_shared_expert_stream 与 dsv4 同名实现近乎逐行重复
modules/dsv4/moe/shared_expert.py:23,43-53已存在_SHARED_EXPERT_STREAM_CACHE与同名函数_ensure_shared_expert_stream(device),函数体与新增的_SHARED_EXPERT_STREAMS/_ensure_shared_expert_stream(generic_moe.py:39-52)近乎逐行相同(同样的非 cuda 早退、同样的device.index or current_device()归一化、同样按 index 缓存)。两处同名不同模块的全局缓存让「本进程每设备到底有几条共享专家侧流」难以推断,也让「每设备一条侧流」这一 capture 前提失去单一权威持有者。需说明:新实现仅在__init__调用(:190),不存在捕获期懒建流路径,故不涉及正确性。 - [6.1] Software Engineering — KISS/YAGNI:无投机性抽象 → issue
四种归约策略靠三个互相纠缠且依赖赋值顺序的布尔标志编码
use_ep_unified_allreduce、use_ep_shared_allreduce、use_unified_tp_allreduce三个布尔量加一个隐式 fallback 共同编码四种互斥策略,且彼此耦合并依赖赋值顺序(use_ep_shared_allreduce必须显式写and not self.use_ep_unified_allreduce,:176)。forward需先用 early return 处理第一种(:311-336),再用if/elif/else处理其余三种(:357-381)。__init__中那段 25 行 ASCII 表格注释(:140-164)本身就是编码已不自解释的证据;再加第五种策略仍需同时改中心谓词块与forward,本次已是第二次这样做。 - [6.1] Software Engineering — LSP:子类/重写保持基类契约 → issue
层未校验 router 返回的即为自己提供的缓冲区,共享专家分支可能被静默丢弃
forward直接对fused_moe的返回值做all_reduce(..., inplace=True)(:336),隐含假设它正是自己传入且已含 gated shared partial 的 buffer;而FusedMoe.forward只断言output.shape == hidden_states.shape(fused_moe.py:357-359),不校验output is row_scatter_target。若某 router 声明支持该能力却在 finalize 中另分配输出,shared 分支贡献被静默丢弃且不报错。契约测试恰好编码了这种不合规行为:_make_fused_moe令router.finalize.return_value = hidden_states.clone()(contract_test:46),而test_forward_passes_row_scatter_target_only_when_supplied只断言键透传、不断言返回同一性,故该文件全绿。 - [6.1] Software Engineering — OCP:本地扩展点优先于修改中心逻辑 → issue
四种归约策略靠三个互相纠缠且依赖赋值顺序的布尔标志编码
use_ep_unified_allreduce、use_ep_shared_allreduce、use_unified_tp_allreduce三个布尔量加一个隐式 fallback 共同编码四种互斥策略,且彼此耦合并依赖赋值顺序(use_ep_shared_allreduce必须显式写and not self.use_ep_unified_allreduce,:176)。forward需先用 early return 处理第一种(:311-336),再用if/elif/else处理其余三种(:357-381)。__init__中那段 25 行 ASCII 表格注释(:140-164)本身就是编码已不自解释的证据;再加第五种策略仍需同时改中心谓词块与forward,本次已是第二次这样做。 - [6.1] Software Engineering — SRP:模块/类职责单一 → issue
四种归约策略靠三个互相纠缠且依赖赋值顺序的布尔标志编码
use_ep_unified_allreduce、use_ep_shared_allreduce、use_unified_tp_allreduce三个布尔量加一个隐式 fallback 共同编码四种互斥策略,且彼此耦合并依赖赋值顺序(use_ep_shared_allreduce必须显式写and not self.use_ep_unified_allreduce,:176)。forward需先用 early return 处理第一种(:311-336),再用if/elif/else处理其余三种(:357-381)。__init__中那段 25 行 ASCII 表格注释(:140-164)本身就是编码已不自解释的证据;再加第五种策略仍需同时改中心谓词块与forward,本次已是第二次这样做。 - [6.1] Tests — 分布式/跨平台变更有对应覆盖 → issue
侧流的真实装配与端到端 fork→join 接线在所有用例中零覆盖
_make_layer的权重恒为 CPU 张量(:58-61),__init__依self.w1.device取流(generic_moe.py:190)而_ensure_shared_expert_stream对非 cuda 返回 None(:45-46),故use_ep_unified_allreduce分支下shared_expert_stream/shared_expert_ready恒为 None,:306 反向钉住assertIsNone(row_scatter_ready);SharedExpertSideStreamCaptureTest又用SimpleNamespace(:370-377)调未绑定方法,绕过__init__与forward。于是两条连接无任何断言:CUDA 权重下__init__确实建出 stream/event;forward确实把非空 event 交给fused_moe。删除 kwarg 会因 :306 下标读取而 KeyError,但硬编码 ` - [6.1] Tests — 新逻辑有聚焦单测 + 相关集成/smoke 测试 → issue
finalize 断言仅校验位置参数个数,未固定参数顺序
_extract_extra_finalize_args只检查len(args) != 5并返回args[4](:20-26)。若FusedMoe.forward中router.finalize(...)的实参顺序被调换(如 topk_weights 与 topk_ids 互换),位置参数个数不变,本文件全部用例仍通过,而生产会把权重当成 id 传给 combine kernel。由于 router 用Mock(spec=FusedMoeDataRouter)(:34),spec只限制属性名、不校验子方法实参绑定,这层保护完全落在该 helper 上。 - [6.1] Tests — 边界 case 覆盖(空、单元素、最大值) → issue
契约测试用 object() 冒充 torch.cuda.Event,且两个场景挤在同一方法
:156 与 :177 用object()充当row_scatter_ready,而FinalizeArgs.row_scatter_ready与FusedMoe.forward均声明其为torch.cuda.Event(fused_moe.py:42、267),测试因此没有钉住类型契约。另外test_forward_passes_row_scatter_target_only_when_supplied(:142-166)把「未提供 target」与「同时提供 target 与 ready」两个场景串在同一方法且未用subTest,第一段断言失败会掩盖第二段;「target 有值但 ready 为 None」这一实际生产分支(fused_moe.py:346 的 false 分支)在两个测试文件中均未被断言。
RTP-LLM Checklist
- [I] 代码质量 — 同一功能用统一工具函数 → issue
all_reduce 的 inplace 用法在两条统一归约路径间不一致
新路径用all_reduce(experts_output, group=Group.TP, inplace=True)(:336),而语义完全对称的use_unified_tp_allreduce路径仍是all_reduce(experts_output, group=Group.TP)(:365)。collective_torch.all_reduce的inplace为 keyword-only 且只影响 symm_mem 快路径(collective_torch.py:694,716):非 inplace 时symm_mem_comm.all_reduce(tensor, out=None)会额外torch.empty_like一份 [T, H](symm_mem.py:167-168)。而 pure-TP 路径的_merge_shared_expert_output本就已就地写入experts_output,加inplace=True同样安全。
Python Static-First Checklist
- [P.F] 语言陷阱 — 循环中 lambda late binding → issue
测试在循环体内定义闭包并捕获循环变量
test_joins_the_producer_before_touching_the_buffer在for tp_rank in (0, TP_SIZE - 1)循环体内定义wait_event(seen)(:132-136),闭包捕获每轮重新绑定的target与waited。当前写法因闭包只在本轮同步调用而行为正确,但这是 late-binding 类缺陷的典型形状:后续若把调用挪出循环(例如收集起来统一断言)会静默指向最后一轮的target。 - [P.G] 测试规范 — mock/fake/stub 不得替代本次声称覆盖的生产边界 → issue
finalize 断言仅校验位置参数个数,未固定参数顺序
_extract_extra_finalize_args只检查len(args) != 5并返回args[4](:20-26)。若FusedMoe.forward中router.finalize(...)的实参顺序被调换(如 topk_weights 与 topk_ids 互换),位置参数个数不变,本文件全部用例仍通过,而生产会把权重当成 id 传给 combine kernel。由于 router 用Mock(spec=FusedMoeDataRouter)(:34),spec只限制属性名、不校验子方法实参绑定,这层保护完全落在该 helper 上。 - [P.G] 测试规范 — 数据驱动测试用 pytest.mark.parametrize → issue
契约测试用 object() 冒充 torch.cuda.Event,且两个场景挤在同一方法
:156 与 :177 用object()充当row_scatter_ready,而FinalizeArgs.row_scatter_ready与FusedMoe.forward均声明其为torch.cuda.Event(fused_moe.py:42、267),测试因此没有钉住类型契约。另外test_forward_passes_row_scatter_target_only_when_supplied(:142-166)把「未提供 target」与「同时提供 target 与 ready」两个场景串在同一方法且未用subTest,第一段断言失败会掩盖第二段;「target 有值但 ready 为 None」这一实际生产分支(fused_moe.py:346 的 false 分支)在两个测试文件中均未被断言。
Strengths
__init__的归约策略对照表注释(generic_moe.py:140-164)把「partial 才能合并、already-complete 分支不能参与」这条非显然不变量画成表格并标注 Group.TP 集合通信次数,读者无需反推谓词;use_ep_shared_allreduce上显式的and not self.use_ep_unified_allreduce(:176)保证互斥,并被test_ep_unified_displaces_the_separate_shared_reduce正面钉住。- 跨流内存生命周期两个方向都覆盖且带解释:
hidden_states.record_stream(side_stream)保护主流分配块在侧流读取期不被回收,shared_partial.record_stream(current)保护侧流分配块在主流读写期不被回收(generic_moe.py:250-262)——这是同类改动最常漏掉的一环。 - 侧流上只放共享专家(
skip_allreduce=True),所有集合通信仍留主流,规避了「NCCL 在 capture 期跨多流」这一已知雷区;routed 链抛异常时补做 join(:327-335),并有test_a_failed_routed_chain_joins_the_shared_branch覆盖。 _finalize_row_scatter把 join 放在 DeepEP combine 之后、写 buffer 之前,既拿到最大重叠又保证可见性,并显式处理「本 rank 切片为空但仍须 join」的尾部 rank(router:311-324)。assert combined_x.size(0) == slice_size(router:319)把「combine 返回的正是 dispatch 切片、不带 padding」这条隐含假设显式钉死(替代旧路径的 padding 兜底),并有test_rejects_a_combine_output_that_is_not_the_dispatched_slice反向验证。deepep_low_latency_row_scatter_test.py:152不用任何 mock,以 1/3/5/8/9 个 token 覆盖整除、ragged 尾部、尾部 rank 空切片,并用torch.cat(routed_slices).shape[0] == original_num_tokens顺带证明切片无重叠无遗漏——这正是整个优化成立的代数前提。- 契约校验前置:
FusedMoe.forward:270-288在router.prepare/experts.execute之前就对「router 不支持却传 target」「传 ready 但无 target」抛ValueError,对应用例还断言两者未被调用,避免半执行状态。 SharedExpertSideStreamCaptureTest用真实torch.cuda.graph验证 fork 在 capture 结束前已闭合、replay 会重算(未闭合的 fork 只在 capture 末尾报错,纯 mock 无法发现),并覆盖 decode 多 batch size 共享同一 graph pool 的实际形态。test_layer_only_passes_keywords_fused_moe_accepts用inspect.signature(FusedMoe.forward)反查MagicMock吞掉的多余 kwargs;_extract_extra_finalize_args要求router.finalize收到恰好 5 个位置参数——两处都是对 mock 失真的主动补偿。_tp_token_slice、_combine_input、_make_fused_moe三处提取同时消除生产侧与测试侧的重复分片/反量化逻辑;新增 py_test 目标的 tags/exec_properties/deps 与同目录pure_tp_router_skip_allreduce_test完全一致,未引入新的 CI 平台约定。
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment
Add this suggestion to a batch that can be applied as a single commit.This suggestion is invalid because no changes were made to the code.Suggestions cannot be applied while the pull request is closed.Suggestions cannot be applied while viewing a subset of changes.Only one suggestion per line can be applied in a batch.Add this suggestion to a batch that can be applied as a single commit.Applying suggestions on deleted lines is not supported.You must change the existing code in this line in order to create a valid suggestion.Outdated suggestions cannot be applied.This suggestion has been applied or marked resolved.Suggestions cannot be applied from pending reviews.Suggestions cannot be applied on multi-line comments.Suggestions cannot be applied while the pull request is queued to merge.Suggestion cannot be applied right now. Please check back later.
No description provided.