Prhub

#52575 [Rust Frontend] Simplify data-parallel size ownership

原始 PR 作者 BugenZhao 合并时间 2026-08-18 17:06 文件变更 19 提交数 4 评论 7 代码增减 +188 / -134

执行摘要

统一 DP 大小归属与校验到连接层,精简服务端冗余配置

PR body 明确说明这是 #51245 的 follow-up,聚焦“简化 deployment-wide data-parallel size ownership”:在 #51178 引入显式 DP rank 路由后,原实现中 server Config 和 AppState 各自持有可被覆盖的 data_parallel_size,容易与 client transport 的实际拓扑漂移。目标是把唯一数据源收进 TransportMode,并让直接使用 engine-core-client 的调用方也获得相同的前置校验,同时补一个 dense-DP 回归测试。

值得精读。这个 PR 是“单一数据源 + 校验前移 + 序列化兼容”三件事同时做好的范例:DP 大小归属清晰后,服务端不再需要 override 路径;serde(rename) 让 Rust 内部重命名不破坏 Python wire 契约;校验下沉到 client 使所有入口共享同一套保护。建议重点看 transport.rsclient.rs 的交互,以及 E2E 测试如何同时覆盖默认与 include_dp=false 两条路径。

讨论亮点

该 PR 的 review 没有实质技术交锋,仅有一条 njhill 的 APPROVE(“Thanks @BugenZhao”)和一条说明 fork 场景下自动 review 被禁用的 claude bot 提示。核心设计决策主要来自 PR body:

This PR focuses on simplifying deployment-wide data-parallel size ownership and adding a dense-DP regression test.
The refactor removes redundant state that could diverge from the client transport.

另从 commit 历史可看出演进:先移动 data_parallel_size 到 Bootstrapped transport 模式,再澄清 engine 上报字段语义,最后把 dense-DP world-size 测试隔离到独立进程,说明作者对测试隔离性有明确意识。

实现拆解

  1. 收敛 DP 大小到 TransportMode:在 rust/src/engine-core-client/src/client.rsTransportMode::Bootstrapped 变体新增 data_parallel_size 字段,HandshakeOwner 的部署级 DP 大小通过 engine_count 派生;新增 TransportMode::data_parallel_size() 统一对外取值。
  2. 校验逻辑下沉到 client 层:新增 TransportMode::validate(),检查 DP 大小非零、不超过 u16::MAX + 1(vLLM 两字节 engine 身份上限),并校验 Bootstrapped 的引擎区间 [engine_start_index, engine_end_index) 不超出部署级 DP 大小;EngineCoreClientConfig::validate() 委托给 transport,EngineCoreClient::connect() 在建立任何 ZMQ socket 前先调用 config.validate(),使直接使用 client 的调用方与 serve 路径获得一致体验。
  3. 删除服务端重复状态rust/src/server/src/config.rs 移除 Config.data_parallel_size 字段及对应校验分支,Config::validate() 改委托 transport_mode.validate()rust/src/server/src/state.rs 移除 AppState.data_parallel_size 字段、with_data_parallel_size()data_parallel_size()/get_world_size 与 gRPC server info 改为通过 EngineCoreClient 暴露的部署级值。
  4. 澄清 Ready 字段语义rust/src/engine-core-client/src/protocol/handshake.rs 将 Rust 字段重命名为 effective_data_parallel_size,并用 #[serde(rename = "data_parallel_size")] 保持 Python msgpack wire key 不变;dense 独立 DP 引擎上报 1,真正的部署级大小由 client transport 拥有。同步更新 mock_engine.rsgrpc/tests.rscli.rscli/tests.rs 及示例文件。
  5. 测试与 CI 配套:新增 tests/v1/distributed/test_dense_dp_world_size.py 的 dense-DP E2E 测试,覆盖默认响应(含 DP)与 include_dp=false;新增 rust/src/engine-core-client/src/tests/client.rsclient_config_validates_bootstrapped_dp_range,验证前端可只拥有全局 DP rank 子集、越界配置被拒绝;rust/src/server/src/routes/tests.rs 删除依赖 override 路径的旧 world-size 测试,并更新 .buildkite/test-amd.yamltest_areas/rust_frontend.yamltest_areas/distributed.yaml 接入新测试。
文件 模块 状态 重要度
rust/src/engine-core-client/src/client.rs 连接层 modified 7.99
rust/src/server/src/config.rs 服务配置 modified 6.39
rust/src/server/src/state.rs 路由状态 modified 5.47
rust/src/engine-core-client/src/protocol/handshake.rs 握手协议 modified 5.28
tests/v1/distributed/test_dense_dp_world_size.py 分布式测试 added 6.25
rust/src/engine-core-client/src/tests/client.rs 客户端测试 modified 5.67
rust/src/server/src/routes/tests.rs 路由测试 modified 7.11

关键符号

TransportMode::data_parallel_size TransportMode::validate EngineCoreClientConfig::validate EngineCoreClient::connect EngineCoreClient::data_parallel_size test_dense_dp_world_size client_config_validates_bootstrapped_dp_range

关键源码片段

rust/src/engine-core-client/src/client.rs core-logic

核心变更文件:Bootstrapped 新增 data_parallel_size 字段,新增 TransportMode::data_parallel_size()/validate(),并在 connect() 前统一执行拓扑校验,是整次重构的枢纽。

impl TransportMode {
    /// 返回 deployment-wide 的数据并行大小。
    ///
    /// HandshakeOwner 模式下 Rust 进程负责全部 engine 的启动协商,`engine_count`
    /// 即为全局 DP 大小;Bootstrapped 模式下 supervisor 可能把全局 DP rank 分割到
    /// 多个 frontend,因此必须显式记录 `data_parallel_size`。
    pub fn data_parallel_size(&self) -> usize {
        match self {
            Self::HandshakeOwner { engine_count, .. } => *engine_count,
            Self::Bootstrapped {
                data_parallel_size, ..
            } => *data_parallel_size,
        }
    }    /// 在打开任何 socket 之前校验 transport 拓扑,避免配置错误延迟到引擎连接阶段
    /// 才暴露。该校验由 `EngineCoreClient::connect` 统一调用,因此直接使用 client
    /// 的调用方与 serve 路径体验一致。
    pub fn validate(&self) -> Result<()> {
        let data_parallel_size = self.data_parallel_size();
        if data_parallel_size == 0 {
            bail_invalid_client_config!("data parallel size must be at least 1");
        }
        if data_parallel_size > usize::from(u16::MAX) + 1 {
            bail_invalid_client_config!(
                "data parallel size ({data_parallel_size}) exceeds the two-byte engine identity limit"
            );
        }        match self {
            Self::HandshakeOwner { .. } => {}
            Self::Bootstrapped {
                engine_start_index,
                engine_count,
                ..
            } => {
                // Bootstrapped 模式允许前端只拥有全局 DP rank 的子集(例如 supervisor
                // 把 4 个 rank 分给两个 frontend 各 2 个),但引擎区间必须落在部署级
                // DP 大小之内。
                if *engine_count == 0 {
                    bail_invalid_client_config!("engine count must be at least 1");
                }
                let engine_start_index = usize::try_from(*engine_start_index).map_err(|_| {
                    Error::InvalidClientConfig {
                        message: "engine start index does not fit usize".to_string(),
                    }
                })?;
                let engine_end_index =
                    engine_start_index.checked_add(*engine_count).ok_or_else(|| {
                        Error::InvalidClientConfig {
                            message: "engine start index + engine count overflows".to_string(),
                        }
                    })?;
                if engine_end_index > data_parallel_size {
                    bail_invalid_client_config!(
                        "connected engine range [{engine_start_index}, {engine_end_index}) exceeds data parallel size ({data_parallel_size})"
                    );
                }
            }
        }        Ok(())
    }
}
rust/src/engine-core-client/src/protocol/handshake.rs data-contract

Ready 响应字段重命名为 effective_data_parallel_size,并用 serde rename 保持 Python wire key 不变,是语义澄清与兼容性的核心。

pub struct EngineCoreReadyResponse {
    // ... 其他握手元数据字段省略 ...    /// 本 EngineCore 有效并行配置中的 data-parallel 大小。
    ///
    /// dense 独立 DP 引擎会被重配置为上报 `1`,deployment-wide 数值由
    /// client transport 的 `TransportMode` 持有,避免引擎上报值与前端
    /// 路由拓扑不一致,也避免不同 frontend 各自持有可漂移的副本。
    #[serde(rename = "data_parallel_size")]
    pub effective_data_parallel_size: u64,
}

评论区精华

deployment-wide DP 大小的单一归属 设计

PR body 明确将 deployment-wide DP size 只存储在 TransportMode::Bootstrapped,HandshakeOwner 从 engine_count 派生;server Config 和 AppState 中的重复字段与 override 路径被整体删除。

结论:按设计合并:Config::validate 改为委托 transport_mode.validate(),/get_world_size 和 gRPC server info 统一从 client 读取部署级值。 · 已解决

engine 上报字段语义澄清与 wire 兼容 设计

Ready 响应字段从 data_parallel_size 重命名为 effective_data_parallel_size,但通过 #[serde(rename="data_parallel_size")] 保持 Python msgpack wire key 不变;dense 独立 DP 引擎上报 1,部署级数值由 client transport 持有。

结论:合并时保留;grpc/tests 中 server info 的 DP 断言相应从 4 改为 2,表明不再使用 frontend override。 · 已解决

风险与影响

  1. Rust 内部 API breakingConfig.data_parallel_sizeAppState::with_data_parallel_size 被删除,任何外部直接构造 vllm-server::Config 的代码(如 external_engine_openai_qwen.rs)都会编译失败,本 PR 已同步更新示例与测试,但仓库外的下游自定义入口需要跟着迁移。
  2. 校验时机前置的行为变化connect() 现在在打开 socket 前拒绝非法拓扑,之前可能延迟到引擎连接阶段才暴露;配置错误的调用方会得到更早更明确的失败,但原本“能启动但最终报错”的场景会变成“直接拒绝启动”,需确认所有生产入口都已传对参数。
  3. 字段语义混淆风险effective_data_parallel_size 对 dense 独立 DP 引擎固定上报 1,若后续代码误把它当作部署级 DP 大小读取,会得到错误值;gRPC server info 的断言从 4 改为 2,说明已不再依赖 frontend override,但任何遗留依赖需彻底清理。
  4. 新 E2E 测试依赖 GPU 分布式 CItest_dense_dp_world_size.py 需要真实拉起 dense-DP 服务,若 CI 分片或 GPU 资源不足可能产生 flaky,且依赖 Qwen/Qwen3-0.6B 模型下载。

影响范围集中在 Rust frontend 内部架构:

  • 系统层面:deployment-wide DP 大小从 Config/AppState 的“可覆盖副本”收敛为 TransportMode 的唯一数据源,消除状态漂移;/get_world_size 与 gRPC server info 的行为由 #51178 引入的语义保持一致。
  • 调用方体验:直接使用 engine-core-client 的第三方获得与 serve 一致的前置拓扑校验,错误配置在 socket 建立前即被拒绝。
  • 团队协作:字段语义澄清(effective vs. deployment-wide)降低了后续维护者误用 engine 上报值的概率;测试从 Rust mock 迁移到真实 dense-DP E2E,回归保障更强。
核心路径变更 跨模块重构 Rust 内部 API breaking 新 E2E 测试依赖 GPU CI

关联 Issue

#51178 [Rust Frontend][gRPC] Add explicit data-parallel rank routing

完整报告

参与讨论