Prhub

#33370 [Feature] Add process-local in-memory KV indexer and Router integration

原始 PR 作者 wuyl1 合并时间 2026-08-20 10:45 文件变更 42 提交数 12 评论 32 代码增减 +6605 / -176

执行摘要

新增进程内 KV 索引器并集成 Router 前缀路由

RFC #31458 提出 KV Indexer 架构,用于跟踪 worker 的 KV 块放置,支持缓存感知路由;#32662 依赖 Redis 等外部存储。本 PR 以更小的里程碑替代 #32662,移除外部数据库,让 event → index → prefix-query 的核心链路更小、更易运行和测试。PR body 说明索引器是元数据软状态服务,KV 数据仍由 worker 持有,并为后续持久化与恢复建立经过验证的单进程基线。

值得精读,尤其适合关注“实验性系统如何通过降级语义换取可运维性”的读者。建议重点阅读 memory_backend.rsRwLock 状态设计、service.rs 中作为前缀语义权威定义的默认实现、client.rs 对错误类别的精细区分,以及 admission.rs 基于调用方 deadline 的脱落。这份 PR 也是后续分布式或持久化 KV Indexer 的基线参考。

讨论亮点

评审聚焦资源边界与数据一致性。HZH0425 指出 2048 块扫描上限会让 Router 用全量请求长度计算匹配率,导致全缓存长前缀被低估并退回 min-load;作者确认后移除上限,改为按 16,384 哈希分块扫描并保持跨块匹配状态,同时补回归测试。另一争议是 gRPC 消息大小:page_size=1 时约 19 万 hashes 即触碰 tonic 默认 4 MiB 上限;作者对比流式分割与打包 sfixed64 编码两种方案,HZH0425 选择打包编码,最终把解码上限提到 8 MiB,约支持 100 万块。此外还确认外部 Indexer 启用后 Router 不再维护本地 KV 树,并补齐全链路 E2E 测试。

实现拆解

变更分 6 步:

  1. 新增 crate 与契约sgl-kv-indexer 作为 Router workspace 成员,共享 lockfile 与工具链。crate 内包含 proto 契约、bridge.rsservice.rsmemory_backend.rsclient.rs 和可执行入口 kv-indexer-server。CI 增加 protoc 安装,保证 proto 编译。
  2. 事件接收与转发(Bridge)bridge.rs 订阅 worker 的 ZMQ KV-event 流,用 rmpv 解码 BlockStored 批次,转换为 ApplyExternalKvBatch 的 action 列表,再通过 gRPC 发给 Indexer。超大批次按 MAX_ACTIONS_PER_BATCHMAX_HASHES_PER_REQUEST 切分为有序 RPC,保持动作顺序与元数据对齐。错误分类区分永久错误与瞬时错误,连接断开以 500 ms–10 s 退避重连,但刻意不做事件重放。
  3. 进程内放置索引(memory_backend)memory_backend.rsRwLock<State> 保存全量放置视图:blocks 记录块 hash 到 (worker, tier) 组件掩码与 token 数,workers 记录 worker 地址与 WorkerCacheSpecholdings 提供按 tier 的反查索引。apply() 是唯一写路径:REPORT 是替换式快照,同批重复 hash 保留最后一次;CLEAR_ALL_AT_TIER 靠反查索引批量清理。读路径通过 collect_worker_prefix_inputs 在一次读锁快照内收集所有候选 worker 的组件放置。
  4. 组件感知前缀匹配(service)service.rs 定义 KvIndexerBackend trait,默认的 match_external_kv_prefix 就是前缀语义的书面定义:逐块应用 FULL(必须连续存在)、SWA(候选边界处覆盖尾部窗口或从头连续)、MAMBA(仅边界块)规则,并返回每个 worker 的最长连续可复用前缀。协议层在进入后端前强制 MAX_HASHES_PER_REQUEST = 16_384MAX_ACTIONS_PER_BATCH = 256,并用 MAX_GRPC_DECODING_MESSAGE_SIZE 限制解码层消息大小。
  5. Router 集成与降级:CLI 新增 --kv-indexer-endpoint--kv-indexer-query-timeout-ms--kv-indexer-query-max-inflight 等参数,并校验必须与 cache_aware_zmq 策略配合。cache_aware_zmq.rsselect_external 把外部前缀信号与当前候选 worker 地址求交,在最优前缀持有者中按最小活动负载选择;无信号(unreachable/timeout/overloaded/query-too-large/empty)一律降级为 min-load 路由并打 WARN,只有被拒(Rejected)才以 503 失败。启用外部 Indexer 后 Router 不再维护本地 radix 树。
  6. 测试与 CI 配套:新增内存后端集成测试、gRPC 契约测试、deadline 脱落测试、CLI 校验测试,以及覆盖真实 worker → ZMQ → Bridge → Indexer → Router 的 E2E 测试(GPU CI 运行)。同时同步更新文档与 Cargo workspace 配置。
文件 模块 状态 重要度
experimental/sgl-router/sgl-kv-indexer/src/bridge.rs 事件桥接 added 9.17
experimental/sgl-router/sgl-kv-indexer/src/memory_backend.rs 存储后端 added 9.17
experimental/sgl-router/sgl-kv-indexer/src/service.rs 索引服务 added 9.28
experimental/sgl-router/sgl-kv-indexer/src/client.rs 查询客户端 added 8.88
experimental/sgl-router/src/policies/cache_aware_zmq.rs 路由策略 modified 8.11
experimental/sgl-router/src/config/cli.rs 配置解析 modified 8.29
experimental/sgl-router/sgl-kv-indexer/src/admission.rs 准入控制 added 8.87
experimental/sgl-router/sgl-kv-indexer/tests/memory_integration.rs 集成测试 added 7.74
experimental/sgl-router/src/server/routes/chat.rs 请求路由 modified 7.82

关键符号

apply_external_kv_batch match_external_kv match_external_kv_prefix collect_worker_prefix_inputs get_external_kv_hit_counts apply do_match_prefix select_external resolve_prefix_query reject_if_deadline_passed stamp_arrival match_prefix from_env

关键源码片段

experimental/sgl-router/sgl-kv-indexer/src/service.rs core-logic

定义 gRPC 服务 trait 与协议边界,默认实现组件感知前缀匹配的权威语义,并强制请求级资源上限。

/// 存储后端 trait。所有变更走 `apply_external_kv_batch`,保持唯一有序写路径。
/// async 允许未来接入 I/O 后端;dyn-safe 让服务端持有 `Arc<dyn KvIndexerBackend>`。
#[tonic::async_trait]
pub trait KvIndexerBackend: Send + Sync + 'static {
    // 应用整个 KVEventBatch。动作预校验,必须按序应用;seq 仅作参考。
    async fn apply_external_kv_batch(
        &self,
        request: ApplyExternalKvBatchRequest,
    ) -> Result<ApplyExternalKvBatchResponse, Status>;    async fn match_external_kv(
        &self,
        request: MatchExternalKvRequest,
    ) -> Result<MatchExternalKvResponse, Status>;    // 收集各 worker 的组件放置数据,与 hashes 对齐。
    // 默认实现是组件盲的:持有块都作为 legacy 整块放置。
    // 组件感知后端应覆盖此方法并附加 WorkerCacheSpec 与驻留组件集。
    async fn collect_worker_prefix_inputs(
        &self,
        hashes: &[i64],
    ) -> Result<Vec<WorkerPrefixInput>, Status> {
        let matched = self
            .match_external_kv(MatchExternalKvRequest {
                hashes: hashes.to_vec(),
                count_as_hit: false,
            })
            .await?;
        Ok(legacy_inputs_from_match(hashes, &matched))
    }    /// 回答每个 worker 持有的最长可复用请求前缀。
    /// 默认实现本身就是前缀语义的“书面定义”;后端若为性能覆盖,必须保持
    /// 逐字段一致(仅允许 `blocks_read` 不同,它是可观测性而非语义)。
    async fn match_external_kv_prefix(
        &self,
        request: MatchExternalKvPrefixRequest,
    ) -> Result<MatchExternalKvPrefixResponse, Status> {
        let limit = prefix_limit(request.hashes.len(), request.max_blocks);
        let hashes: Vec<i64> = request.hashes.into_iter().take(limit).collect();
        if hashes.is_empty() {
            return Ok(MatchExternalKvPrefixResponse::default());
        }
        let inputs = self.collect_worker_prefix_inputs(&hashes).await?;
        Ok(compute_prefix_response(&inputs, hashes.len() as u32))
    }
}
experimental/sgl-router/sgl-kv-indexer/src/admission.rs core-logic

查询路径过载保护:基于调用方 deadline 做脱落,并用翻倍计数抑制日志风暴。

//! 基于调用方截止时间的查询路径负载脱落。
//! 一个等待超过调用方完整截止时间的查询已经无法被有价值地回答,
//! 服务它只会拖延其余积压。预算来自调用方自己的 `grpc-timeout`;
//! 没有服务器侧阈值需要调,未声明截止时间的调用方永不被脱落。use std::sync::atomic::{AtomicU64, Ordering};
use std::time::{Duration, Instant};
use tonic::metadata::MetadataMap;
use tonic::{Extensions, Request, Status};/// 请求到达时刻,在请求头被读取之后打上。
#[derive(Clone, Copy)]
struct Arrival(Instant);/// 按类型统计拒绝次数,并在总数翻倍时上报,避免日志风暴。
pub(crate) struct RejectionLog(AtomicU64);impl RejectionLog {
    pub(crate) const fn new() -> Self { Self(AtomicU64::new(0)) }    // 记录一次拒绝;在总数达到 2 的幂时返回总数供日志使用。
    pub(crate) fn record(&self) -> Option<u64> {
        let total = self.0.fetch_add(1, Ordering::Relaxed) + 1;
        total.is_power_of_two().then_some(total)
    }
}static DEADLINE_SHED_LOG: RejectionLog = RejectionLog::new();// 在请求被调度前打上到达时间戳;没有该拦截器就不会发生脱落。
pub fn stamp_arrival(mut request: Request<()>) -> Result<Request<()>, Status> {
    request.extensions_mut().insert(Arrival(Instant::now()));
    Ok(request)
}// 对耗时超过调用方声明 deadline 的查询返回脱落。
// 绝不用于 apply 路径:丢弃 KV 事件会让索引与 worker 报告永久分叉。
pub(crate) fn reject_if_deadline_passed(
    metadata: &MetadataMap,
    extensions: &Extensions,
) -> Result<(), Status> {
    let (Some(arrival), Some(budget)) = (extensions.get::<Arrival>(), caller_deadline(metadata))
    else {
        return Ok(());
    };
    let waited = arrival.0.elapsed();
    if waited < budget {
        return Ok(());
    }
    if let Some(shed_total) = DEADLINE_SHED_LOG.record() {
        tracing::info!(
            shed_total,
            waited_ms = waited.as_millis(),
            budget_ms = budget.as_millis(),
            "shedding prefix query whose caller deadline already passed"
        );
    }
    Err(Status::deadline_exceeded(
        "prefix query waited longer than its caller deadline",
    ))
}

评论区精华

2048 扫描上限与匹配率误差 正确性

hzh0425 问为何限制 2048 块,并指出 Router 用全量请求长度计算匹配率,未扫描块会被当作 miss,导致 4096 块的完全缓存请求显示成 50% 命中而回退到 min-load 路由。

结论:wuyl1 确认问题,移除 2048 上限,改为按 16,384 哈希分块扫描并跨块保持匹配状态,同时匹配率按实际扫描范围计算,添加大于 16,384 块的回归测试;commit 8e5fcd4054 修复。 · 已解决

超大批次超过 16,384 hashes 被拒绝 设计

hzh0425 指出 producer 单批次可超 16,384 hashes,而 indexer 会拒绝。wuyl1 最初建议保留服务端限制并让 Bridge 拆分超大批次。

结论:Bridge 拆分超大批次到有序 RPC,保持动作与元数据对齐;已在 commit 8e5fcd4054 修复。 · 已解决

CI 缺少 PROTOC 环境 infra

hzh0425 提醒检查 CI 是否有 PROTOC。wuyl1 确认 Docker 有 protoc 但 lint/build/test CI 没有,计划加入 install_protoc.sh。

结论:Router CI 已安装 protoc。 · 已解决

外部索引器启用时是否维护本地 KV 树 设计

hzh0425 问 external indexer 启用时本地树是否仍维护。wuyl1 反问偏好,hzh0425 回答 'yes',要求停止维护,并同时要求改进策略和增加 e2e 测试。

结论:使用外部 Indexer 时 Router 不再维护本地 KV 树;E2E 测试已在 commit d41c99e4662 添加。 · 已解决

gRPC 消息大小与 page_size=1 性能

hzh0425 指出 page_size=1 时 payload 过大可能导致 gRPC 失败,问是否假定 page_size=64。wuyl1 说不假定;对比流式分割与打包 sfixed64 编码两种方案。

结论:哈希改为打包 `sfixed64`,索引器解码上限调至 8 MiB,约支持一百万个块;commit 39b5bdeaaa 实现并补充回归测试。 · 已解决

E2E 覆盖全链路 测试

hzh0425 要求完整 E2E:真实 worker → ZMQ → Bridge → Indexer gRPC → Router 前缀查询 → 路由到缓存 worker。

结论:wuyl1 在 commit d41c99e4662 增加外部 KV Indexer E2E 覆盖,验证 Worker → ZMQ → Bridge → Indexer → Router 全链路。 · 已解决

风险与影响

  • 软状态与可用性:Indexer 进程重启即丢失全部放置元数据;Bridge 断连期间产生的 KV 事件不会恢复,缓存亲和度会暂时退回到无信号状态。这是有意的取舍,但部署方必须接受。
  • 单进程 RwLock 瓶颈:所有 apply 串行化写锁、查询共享读锁,在 worker 数量或事件频率增大后可能成为热点,尤其是长提示词大量命中场景。
  • gRPC 消息与内存安全:8 MiB 解码上限解决了 page_size=1 场景,但 MAX_CONCURRENT_STREAMS=64 只约束并发流数量,解码发生在后端逻辑之前,恶意或有缺陷的调用方仍可让服务器缓冲 8 MiB × 64 的字节。
  • Router 行为变化:启用外部 Indexer 后停止维护本地 radix 树;若配置的 endpoint 无法拨号,Router 启动即失败(而非运行时逐请求失败)。
  • 构建依赖:proto 编译需要 protoc,已修复 CI,但本地开发环境需要额外工具链。
  • 对用户:默认不启用,不影响现有推理路径;启用后从本地 KV 树缓存感知切换为跨 worker 的外部前缀信号,缓存命中率与长前缀复用是主要收益。
  • 对系统:新增一个独立 Indexer 进程与每 worker 一个 Bridge 进程;Router 增加最多 32 个并发前缀查询(默认),每条查询默认 100 ms deadline。索引器内存占用与跟踪的块数成正比。
  • 对团队:新增 Rust crate 与 20+ 测试文件;Router CI 依赖 protoc;sgl-kv-indexer 与 Router 共享窄特性集 zeromq,扩大构建面。
软状态无持久化 无事件重放 单进程 RwLock 瓶颈 gRPC 8 MiB 解码上限 需要 protoc 构建依赖

关联 Issue

#32514 feat(kv-events): Add component_types field to BlockStored for per-component placement tracking
#32537 fix(hicache): report recomputed GPU blocks

完整报告

参与讨论