Prhub

#32900 [Distributed] Propagate semantic group names to PyTorch process groups

原始 PR 作者 anranxia 合并时间 2026-08-07 11:38 文件变更 2 提交数 7 评论 15 代码增减 +78 / -1

执行摘要

为进程组新增 group_desc 语义元数据,供诊断工具区分 TP/PP

PR body 明确指出:GroupCoordinator 已经携带逻辑 group_name(tp、pp、moe_ep 等),但从未传播到其创建的 ProcessGroup,导致该信息对标准 PyTorch 元数据以及任何检查进程组的 profiler / 诊断工具(如区分 TP 通信器与 PP 通信器)不可见。测试 docstring 进一步说明,这是 NCCL Inspector 区分 TP/PP 通信工具的依据,丢失该元数据会重新引入 PP 被误分类为 TP 的问题。

值得快速精读,尤其关注测试设计:从 mock 断言到真实 new_group 读回的演进,体现了“用实现验证实现”的同义反复陷阱以及回归守卫的正确写法。代码改动本身很小(+7/-1),但作为分布式可观测性元数据约定({group_name}:device/cpu),对后续 NCCL Inspector 等工具链有实际价值。

讨论亮点

review 核心交锋集中在测试有效性上。hnyls2002 指出:

These three tests assert f"{group_name}:device" against the same expression the patch writes, with torch.distributed.new_group itself mocked out ... so they pass regardless of whether new_group actually accepts the kwarg, which is this PR's only real risk.

同时质疑测试文件从 register_cuda_ci + register_amd_ci 迁到 register_cpu_ci 会丢失 GPU/AMD 覆盖。anranxia 承认“mock-tautologies”问题,将测试重写为真实 GroupCoordinator over single-rank gloo world 并读回 ProcessGroup 的 group_desc,验证结果:带变更时返回 ('tp:device','tp:cpu') 等,去掉 group_desc 参数后返回 ('undefined','undefined') 且测试失败,证明是真正的回归守卫;register_cpu_ci 迁移说明为来自上游 #33654 的合并结果。另外,首轮 CI 失败根因是测试 fake 了 mooncake.ep 而代码 import 的是 mooncake.pg,在装有真实 mooncake 的 1-gpu-5090 runner 上 fall through 到真实 loader 报 ModuleNotFoundError,修复方式是 fake 正确子模块并设置 path = []。最终 hnyls2002 APPROVED。

实现拆解

  1. 变更入口:唯一源码改动位于 python/sglang/srt/distributed/parallel_state.py 的 GroupCoordinator.init,共四处 torch.distributed.new_group 调用追加 group_desc 关键字参数,无其他逻辑改动。
  2. 核心改造:mooncake 分支的 device_group(backend="mooncake")与 cpu_group(backend="mooncake-cpu")分别传入 group_desc=f"{group_name}:device" 和 f"{group_name}:cpu";普通分支的 device_group(backend=torch_distributed_backend)同样传 device 标签,gloo cpu_group 传 cpu 标签。group_name 为 None 时沿用既有 "anonymous" 归一化逻辑,产出 anonymous:device / anonymous:cpu,避免出现 None:device。group_desc 是 PyTorch 2.4 起的公共参数,SGLang 固定 PyTorch 2.11,因此无需版本门控。
  3. 测试配套:test/registered/unit/distributed/test_parallel_state.py 新增 _read_group_descs 辅助函数,在单 rank gloo + HashStore 环境下构建真实 GroupCoordinator,再从存活的 ProcessGroup 对象读回 group_desc;新增 test_group_desc_propagated_via_real_new_group(参数化 tp/pp)和 test_group_desc_none_normalized_to_anonymous 两个回归测试,并在 main 入口挂载,保证无 GPU/NCCL 环境下可运行。
  4. 演进过程:共 7 个 commit,含多次 merge main。最初版本是 mock torch.distributed.new_group 的测试,经 review 批评为“同义反复”后,在 a5e914a 重写为真实 new_group 读回验证;期间还修复了测试中 fake mooncake 子模块名错误(mooncake.ep → mooncake.pg)。测试文件从 register_cuda_ci/register_amd_ci 迁移到 register_cpu_ci 并非本 PR 意图,而是合并上游 #33654 带入的变更。
文件 模块 状态 重要度
python/sglang/srt/distributed/parallel_state.py 并行状态 modified 5.27
test/registered/unit/distributed/test_parallel_state.py 并行单测 modified 6.13

关键符号

GroupCoordinator.__init__ _read_group_descs test_group_desc_propagated_via_real_new_group test_group_desc_none_normalized_to_anonymous

关键源码片段

python/sglang/srt/distributed/parallel_state.py core-logic

唯一源码变更文件:GroupCoordinator.__init__ 中四处 torch.distributed.new_group 调用追加 group_desc 元数据,覆盖普通 NCCL/Gloo 与 mooncake device/cpu 全部分支,是所有分布式通信组创建的必经路径。

# GroupCoordinator.__init__ 中创建 device_group 与 cpu_group 的两条分支
# 关键变更:为每个 new_group 调用附加 group_desc 元数据,供 PyTorch profiler /
# NCCL Inspector 等诊断工具区分 TP / PP / EP 等语义通信组,避免 PP 被误判为 TP。
if self.rank in ranks:
    # mooncake 后端分支:device 与 cpu 组分别携带各自的 pg_options,
    # 同时附加 "{group_name}:device" / "{group_name}:cpu" 标签
    device_group = torch.distributed.new_group(
        ranks,
        backend="mooncake",
        pg_options=dev_opts,
        timeout=subgroup_timeout,
        group_desc=f"{group_name}:device",
    )
    cpu_group = torch.distributed.new_group(
        ranks,
        backend="mooncake-cpu",
        pg_options=cpu_opts,
        timeout=subgroup_timeout,
        group_desc=f"{group_name}:cpu",
    )
else:
    # 普通分支:device 组使用指定后端,cpu 组固定 gloo 用于跨进程 CPU 协调,
    # 两者同样携带语义标签;group_name 为 None 时由上层归一化为 "anonymous"
    active_ranks = torch.ones(len(ranks), dtype=torch.int32, device=self.device)
    active_ranks_cpu = torch.ones(len(ranks), dtype=torch.int32)
    pg_options = get_torch_distributed_pg_options(group_name)
    device_group = torch.distributed.new_group(
        ranks,
        backend=torch_distributed_backend,
        pg_options=pg_options,
        timeout=subgroup_timeout,
        group_desc=f"{group_name}:device",
    )
    cpu_group = torch.distributed.new_group(
        ranks,
        backend="gloo",
        timeout=gloo_timeout,
        group_desc=f"{group_name}:cpu",
    )
test/registered/unit/distributed/test_parallel_state.py test-coverage

新增回归测试并经历重大改写:从 mock new_group 的同义反复断言演进为真实 gloo world_size=1 环境读回 ProcessGroup.group_desc,直接守卫本 PR 唯一真实风险(PyTorch 是否接受 group_desc 参数),是评审聚焦的核心文件。

def _read_group_descs(group_name):
    """在单 rank 的 gloo 世界中构建真实 GroupCoordinator,并读回其创建的
    ProcessGroup 上的 group_desc 元数据,返回 (device_group.group_desc,
    cpu_group.group_desc)。    刻意不 mock torch.distributed.new_group,而是驱动真实实现,因此:
    * 若安装的 PyTorch 不接受 group_desc 关键字参数(本变更唯一真实风险),
      测试会在构造时直接失败;
    * 断言的是从 ProcessGroup 读回的值,而不是传给 new_group 的字符串,
      避免“用实现验证实现”的同义反复(缺省时 torch 返回 "undefined")。
    gloo + world_size=1 无需 GPU / NCCL,可保留在 CPU 套件中。
    """
    import torch.distributed as dist
​
    dist.init_process_group(
        backend="gloo", store=dist.HashStore(), rank=0, world_size=1
    )
    try:
        coord = parallel_state.GroupCoordinator(
            group_ranks=[[0]],
            local_rank=0,
            torch_distributed_backend="gloo",
            use_pynccl=False,
            use_pymscclpp=False,
            use_custom_allreduce=False,
            use_torch_symm_mem_all_reduce=False,
            use_hpu_communicator=False,
            use_xpu_communicator=False,
            use_npu_communicator=False,
            use_message_queue_broadcaster=False,
            group_name=group_name,
        )
        return coord.device_group.group_desc, coord.cpu_group.group_desc
    finally:
        if dist.is_initialized():
            dist.destroy_process_group()
​
​
@pytest.mark.parametrize("group_name", ["tp", "pp"])
def test_group_desc_propagated_via_real_new_group(group_name):
    """回归守卫:GroupCoordinator 必须以 group_desc=f"{group_name}:device|cpu"
    标记其创建的进程组。这是 NCCL Inspector 区分 TP / PP 通信元的依据,
    丢失该元数据会重新引入“PP 被误分类为 TP”的问题。"""
    assert _read_group_descs(group_name) == (
        f"{group_name}:device",
        f"{group_name}:cpu",
    )
​
​
def test_group_desc_none_normalized_to_anonymous():
    """group_name 为 None 时保持既有 "anonymous" 归一化,
    生成 anonymous:device / anonymous:cpu 而不是 None:device。"""
    assert _read_group_descs(None) == ("anonymous:device", "anonymous:cpu")

评论区精华

mock 测试是同义反复,无法守卫 group_desc 真实风险 测试

hnyls2002 指出三个测试断言的是 patch 本身写入的同一个 f-string,且 torch.distributed.new_group 被 mock 掉,因此无论 new_group 是否真正接受 group_desc 关键字参数,测试都会通过——而“是否接受该参数”正是本 PR 唯一真实风险。

结论:anranxia 承认问题并重写测试:构建真实 GroupCoordinator over single-rank gloo 世界,从存活 ProcessGroup 读回 group_desc 断言;去除 group_desc 参数后测试返回 'undefined' 并失败,证明是真正的回归守卫。 · 已解决

mooncake.pg fake 子模块导致 CI 失败 正确性

首轮 1-gpu-5090 CI 失败:测试 fake 了 mooncake.ep,但 GroupCoordinator mooncake 分支实际 import 自 mooncake.pg;runner 装有真实(不同构建的)mooncake 时 fall through 到真实 loader,报 ModuleNotFoundError: No module named 'mooncake.pg'。

结论:修复为 fake 正确的子模块 mooncake.pg 并设置 fake_mooncake.__path__ = [],保证 import 解析不受环境真实 mooncake 影响;功能代码无需改动。 · 已解决

测试从 GPU/AMD CI 迁移到 CPU CI 是否丢失覆盖 测试

hnyls2002 质疑文件从 register_cuda_ci + register_amd_ci 移到 register_cpu_ci,可能丢掉两个既有 group-construction 测试的 GPU/AMD 覆盖,并询问运行时变更能否仅靠现有 TP/PP 套件覆盖。

结论:anranxia 说明 register_cpu_ci 迁移并非本 PR 意图,而是合并上游 #33654 带入的变更;本 PR 保留 CPU 套件可运行的真实回归测试即可。 · 已解决

风险与影响

  1. PyTorch 版本兼容:group_desc 参数要求 PyTorch >= 2.4,SGLang 固定 2.11 所以受控,但若运行环境使用更老版本会在进程组创建时抛 TypeError——这也是 PR 自述的唯一真实风险,测试 docstring 已明确将其作为回归守卫目标。
  2. 核心路径变更:GroupCoordinator 是所有 TP/PP/EP/DP 及 mooncake 通信组创建的必经路径,改动虽只是附加 kwarg,但影响面覆盖全部分布式初始化;好在不改变组成员、后端与 collectives,行为回归概率低。
  3. 测试覆盖迁移:测试文件从 GPU/AMD CI 迁移到 CPU CI(上游 #33654 带入),GPU/NCCL 真实路径的覆盖依赖现有 TP/PP 集成测试,mooncake 分支也未被 CPU 单测直接覆盖(作者声明在真实 mooncake 环境双向验证过)。
  4. 无安全/性能风险:group_desc 仅在创建时记录字符串元数据,不影响运行时通信性能。

对用户:无 API 变化、无行为变化,进程组创建时额外携带语义标签,属于纯增量元数据。对系统:所有基于 GroupCoordinator 的分布式场景(TP、PP、EP、DP、mooncake 等)创建的进程组都会带上 {group_name}:device/cpu 描述,PyTorch profiler、NCCL Inspector 等诊断工具可据此区分通信组语义,消除 PP 被误判为 TP 的困惑。对团队:为分布式可观测性基础设施打下基础,后续诊断、监控、告警工具可直接依赖该元数据做更精确的通信拓扑分析。

所有进程组创建核心路径 依赖 PyTorch >= 2.4 的 group_desc 测试迁移至 CPU CI,GPU 覆盖降低 mooncake 分支未纳入 CPU 单测

关联 Issue

未识别关联 Issue

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

完整报告

参与讨论