Prhub

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

原始 PR 作者 connorcarpenter15 合并时间 2026-08-11 06:05 文件变更 19 提交数 8 评论 13 代码增减 +215 / -60

执行摘要

为 Rust gRPC 前端新增显式 DP rank 引擎路由

PR body 的 Purpose 明确说明动机:为 Rust frontend gRPC 推理 API 增加显式 data-parallel rank 路由,并“Reject ranks that cannot be represented by vLLM's two-byte engine identity instead of allowing integer truncation to alias another rank”。在分体推理、KV 亲和路由等场景下,调用方需要把请求定向投递到指定 DP 副本,而 Rust 前端此前只支持本地索引 + 负载均衡,无法表达全局 rank 语义。njhill 在评论中也强调 rank 应作为 gRPC metadata 而非 proto 字段承载(引用早期 PR #48033 的讨论)。

值得精读。核心看 choose_engine_for_request 的过滤式写法、metadata 解析的 fail-loud 策略、config.rs 的校验矩阵。review 中 njhill 对 DP size 归属的判断、BugenZhao 对能力字段的反对,体现了“少加协议字段、保持前端职责单一”的设计哲学,对后续 Rust 前端协议演进有直接借鉴价值。

讨论亮点
  1. DP size 应放在哪一层:BugenZhao 认为理论上应通过握手 EngineCoreReadyResponse 传给前端,但 dense 模型会覆盖 ParallelConfig.data_parallel_size,因此开了 #51245 修复;njhill 则认为更稳妥的是沿用作者最初的方案——“前端保留外部配置的 DP size、引擎按 DP=1 工作”,与 Python 前端行为一致。最终保留 frontend 持有的 data_parallel_size 配置字段,后续再随 #51245 迁移。
  2. rank 走 metadata 还是 proto 字段:njhill 在合并前一直坚持 gRPC metadata 方案,作者通过 commit refactor(grpc): route data parallel rank through metadata 最终采纳。
  3. 能力字段取舍:BugenZhao 反对新增 supports_explicit_data_parallel_rank 布尔字段,认为早期阶段无需维护向后兼容,可用 engine_version 检测,否则每个协议演进都要加 flag;作者同意并移除。
  4. u16 收紧是否正交:BugenZhao 询问 from_engine_index(u32 -> u16) 是否属于正交变更;作者解释把校验收敛到边界更合理,移除也不影响功能。
  5. 错误信息增强:BugenZhao 评价 connected_ranks 是 “good catch”,是添加 hybrid/external DPLB 时遗漏的问题。
  6. njhill 最终 LGTM:建议把 configured_data_parallel_size 改名为 data_parallel_size,并删除 client.rs 中不再使用的函数。

实现拆解

本次变更按以下 4 步落地:

  1. gRPC 入口接收 rank(rust/src/server/src/grpc/inference.rs:新增 metadata 常量 DATA_PARALLEL_RANK_METADATA_KEY = "x-data-parallel-rank" 与解析函数 data_parallel_rank_from_metadata()generategenerate_stream 两个 RPC 都在 prepare_request 前解析并注入,最终写入 TextRequest.data_parallel_rank。选择 metadata 而非 proto 字段是采纳 njhill 在早期 PR #48033 中的建议,避免路由信息入侵请求 schema。

  2. 核心路由按全局身份匹配(rust/src/engine-core-client/src/client/state.rs + transport.rschoose_engine_for_request 在指定 rank 时先 u16::try_from(rank),失败即返回 InvalidDataParallelRank,成功后再校验该引擎是否已连接,错误信息携带 connected_ranks 全量列表。配套把 EngineId::from_engine_index 签名从 u32 收紧为 u16apply_scheduler_counts 增加 u16 防护,connect_bootstrappedengine_start_index + engine_count 做区间边界检查,整个编码链路不再出现静默截断。

  3. 部署级 DP 拓扑持有与校验(rust/src/server/src/config.rs + state.rs + control.rs/world_size.rsConfig 新增 data_parallel_size 字段,validate() 校验其 ≥1、≤ u16::MAX + 1,并与 HandshakeOwner(要求 engine_count 相等)或 Bootstrapped(要求引擎范围不越界)的引擎配置交叉验证;AppState 新增该字段及 with_data_parallel_size() 注入方法;server_info 与 world_size 路由改用该配置值上报,不再读引擎握手值(因此删除了 EngineCoreClient::data_parallel_size() 死代码)。Python 侧 vllm/entrypoints/cli/serve.pyvllm/v1/utils.py 新增 --data-parallel-size 参数透传给 Rust 前端。

  4. 配套收尾rust/src/mock-engine/src/lib.rsrun_engine 改为 usize 参数并在内部做 u16::try_from 校验;InvalidDataParallelRank 错误类型重构为携带 connected_ranks: Vec<u32>;测试以内嵌 Rust 单元 / 集成测试为主(state.rs 测试模块、grpc/tests.rs),无独立测试文件:新增 unary_generate_unconnected_data_parallel_rank_returns_invalid_argument(全局 rank 3 的引擎拒绝 rank 0 请求)、register_with_rank_uses_global_engine_identity(验证全局身份语义)、control_aggregates_multi_engine_capacity 扩展(server_info 上报 data_parallel_size)。

文件 模块 状态 重要度
rust/src/engine-core-client/src/client/state.rs 请求路由 modified 7.78
rust/src/server/src/grpc/inference.rs gRPC 服务 modified 6.92
rust/src/engine-core-client/src/transport.rs 传输层 modified 7.12
rust/src/server/src/config.rs 配置 modified 6.38
rust/src/server/src/state.rs 应用状态 modified 5.87
rust/src/server/src/grpc/tests.rs 服务测试 modified 6.38
rust/src/engine-core-client/src/error.rs 错误处理 modified 5.52
vllm/v1/utils.py 启动编排 modified 4.99
rust/src/mock-engine/src/lib.rs 模拟引擎 modified 6.02
vllm/entrypoints/cli/serve.py CLI 入口 modified 4.35

关键符号

choose_engine_for_request data_parallel_rank_from_metadata from_engine_index connect_bootstrapped with_data_parallel_size data_parallel_size run_engine apply_scheduler_counts

关键源码片段

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

核心路由逻辑所在:`choose_engine_for_request` 实现全局 DP rank 精确匹配与 `u16::try_from` 防截断,错误通过 `connected_ranks` 列表报告;`apply_scheduler_counts` 增加 u16 防护。

/// 为请求选择目标引擎:未指定 rank 时按负载均衡,指定时按全局
/// data-parallel rank 精确匹配已连接引擎的身份。
fn choose_engine_for_request(&mut self, data_parallel_rank: Option<u32>) -> Result<EngineId> {
    if let Some(rank) = data_parallel_rank {
        // vLLM EngineCore 的 ROUTER/DEALER 身份是两字节小端编码,
        // 先把 u32 rank 无损收敛到 u16;否则直接视为无效 rank,
        // 防止整数截断后路由到身份编码恰好相同的另一个引擎。
        let engine_id = u16::try_from(rank).ok().map(EngineId::from_engine_index);
        // 同时要求“编码成功”且“该引擎确实连接到了本前端”,
        // 二者缺一都返回 InvalidDataParallelRank。
        return engine_id
            .filter(|engine_id| self.routing_per_engine.contains_key(engine_id))
            .ok_or_else(|| Error::InvalidDataParallelRank {
                rank,
                // 错误信息带上已连接 rank 的全集,调用方可以区分
                // “rank 写错”和“目标引擎还没连上”两种情况。
                connected_ranks: self
                    .routing_per_engine
                    .keys()
                    .filter_map(EngineId::engine_index)
                    .collect(),
            });
    }    // 未指定 rank:沿用负载均衡,选择路由得分最小的引擎。
    Ok(self
        .routing_per_engine
        .iter()
        .min_by_key(|(_, state)| state.routing_score())
        .map(|(engine_id, _)| engine_id.clone())
        .expect("request registry must contain at least one engine"))
}
rust/src/server/src/grpc/inference.rs core-logic

gRPC 入口:新增 `x-data-parallel-rank` metadata 常量与 `data_parallel_rank_from_metadata()`,在 generate / generate_stream 两个 RPC 中解析并注入 `TextRequest.data_parallel_rank`。

/// gRPC metadata 键名:请求显式指定 data-parallel rank。
///
/// rank 属于“如何路由这次调用”的传输层信息,刻意不放进 `GenerateRequest`
/// proto 结构体,避免路由语义入侵请求 schema(njhill 在早期 PR #48033
/// 中坚持的设计)。
const DATA_PARALLEL_RANK_METADATA_KEY: &str = "x-data-parallel-rank";/// 解析 metadata 中的可选 rank:缺省返回 `None`(走负载均衡),
/// 值必须是可解析为 u32 的字符串,否则以 `InvalidArgument` 拒绝,
/// 避免非法输入被静默忽略。
fn data_parallel_rank_from_metadata(
    request: &Request<pb::GenerateRequest>,
) -> Result<Option<u32>, Status> {
    let Some(value) = request.metadata().get(DATA_PARALLEL_RANK_METADATA_KEY) else {
        return Ok(None);
    };
    let value = value.to_str().map_err(|_| {
        Status::invalid_argument("x-data-parallel-rank metadata must be an unsigned 32-bit integer")
    })?;
    value.trim().parse::<u32>().map(Some).map_err(|_| {
        Status::invalid_argument("x-data-parallel-rank metadata must be an unsigned 32-bit integer")
    })
}
rust/src/engine-core-client/src/transport.rs core-logic

`EngineId::from_engine_index` 从 u32 收紧为 u16,类型层面杜绝截断;`connect_bootstrapped` 对 engine_start_index + engine_count 增加两字节身份区间校验。

if engine_count == 0 {
    bail_unexpected_handshake_message!("expected engine_count >= 1");
}
// bootstrapped 模式的引擎身份由 [engine_start_index, +engine_count)
// 连续合成,必须先验证整个区间落在两字节引擎身份范围内,否则
// EngineId 编码会静默回绕,导致路由别名。
let engine_start_index =
    u16::try_from(engine_start_index).map_err(|_| Error::UnexpectedHandshakeMessage {
        message: "engine_start_index exceeds the two-byte engine identity limit".to_string(),
    })?;
let engine_end_index =
    usize::from(engine_start_index).checked_add(engine_count).ok_or_else(|| {
        Error::UnexpectedHandshakeMessage {
            message: "engine_start_index + engine_count overflows".to_string(),
        }
    })?;
if engine_end_index > usize::from(u16::MAX) + 1 {
    return Err(Error::UnexpectedHandshakeMessage {
        message: "engine_start_index + engine_count exceeds the two-byte engine identity limit"
            .to_string(),
    });
}

评论区精华

DP size 应放在哪一层:前端配置 vs 握手传递 设计

BugenZhao 提出理论上应从 `EngineCoreReadyResponse` 握手获取真实 DP size,但 dense 模型下 `ParallelConfig.data_parallel_size` 会被覆盖,因此开了 #51245;njhill 认为更稳妥的是沿用作者最初方案——前端保留外部配置值、引擎按 DP=1 工作,与 Python 前端行为一致;作者确认前端是唯一消费方。

结论:保留 frontend 持有的 `data_parallel_size` 配置字段,后续随 #51245 落地再迁移。 · 已解决

rank 用 gRPC metadata 而非 proto 字段 设计

njhill 在 issue 评论中坚持 metadata 方案,引用早期 PR #48033 的评论;作者通过 commit `refactor(grpc): route data parallel rank through metadata` 最终采纳。

结论:采纳 metadata 方案,proto 不受路由信息入侵。 · 已解决

supports_explicit_data_parallel_rank 能力字段是否必要 设计

BugenZhao 认为早期阶段无需维护向后兼容,可用 `engine_version` 检测,否则每个协议演进都要加 flag,会迅速膨胀;作者同意并移除该字段。

结论:移除 `supports_explicit_data_parallel_rank` 字段。 · 已解决

InvalidDataParallelRank 错误信息增强 正确性

BugenZhao 评价 `connected_ranks` 是 “good catch”,是添加 hybrid/external DPLB 时遗漏的问题。

结论:错误信息采用已连接 rank 全集,提升可操作性。 · 已解决

from_engine_index 收紧 u16 是否正交变更 question

BugenZhao 询问 `u32 -> u16` 签名收紧是否属于正交变更;作者解释把校验收敛到边界更合理,移除也不影响功能。

结论:保留收紧,作为防截断的一致性改动。 · 已解决

风险与影响

  • metadata 解析 fail-louddata_parallel_rank_from_metadata 对非 u32 字符串直接返回 InvalidArgument,调用方升级后一旦带上非法值会立即收到新错误路径,需要客户端感知。
  • 类型收紧联动面EngineId::from_engine_indexu16 波及 mock-engine、测试与调度统计路径;apply_scheduler_counts 对超范围索引静默忽略(返回 false),虽然现实部署远达不到 65536,但存在监控盲区。
  • 配置一致性假设config.rsHandshakeOwner 模式强制 engine_count == data_parallel_size,对部分引擎连到本前端的混合部署形态可能直接拒绝启动,需与部署团队确认。
  • 协议无前后兼容承诺:BugenZhao 明确建议早期阶段不维护滚动升级兼容,客户端与 server 需同步升级。
  • 调用方(客户端):获得确定性路由能力,可将请求精确投递到指定 DP 副本,是分体推理、KV 亲和路由等场景的硬需求;未设置 x-data-parallel-rank 时行为完全不变。
  • 系统:前端路由语义从“本地索引”升级为“全局身份”,与 Python 前端 DP 语义对齐;错误报告从 num_engines 变为 connected_ranks,可观测性提升。
  • 团队与后续:Rust 前端与 Python CLI 两条链路均需改动;该 PR 与 #51245(握手传递 DP size)存在承接关系,部署拓扑信息来源未来可能迁移。
核心路由路径变更 跨语言配置透传 类型收紧联动面广 协议无前后兼容承诺

关联 Issue

未识别关联 Issue

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

完整报告

参与讨论