Prhub

#47444 [Rust Frontend] Cache metric handles for scheduler & request stats

原始 PR 作者 BugenZhao 合并时间 2026-07-06 21:03 文件变更 9 提交数 4 评论 5 代码增减 +343 / -237

执行摘要

缓存 scheduler 和 request 级 metrics handles 以减少重复查找

根据 PR body:'This PR reduces repeated get_or_create lookups on scheduler & request-level metrics paths. Instead of calling get_or_create every time when we want to access the global metrics handle for a request, we pre-resolve per-engine metrics handles once ClientInner and GenerateOutputStream are created, and cache them for future use.'

值得精读,展示了 Rust 中通过预解析和缓存 metrics handles 来优化性能的实践模式。适合作为 Rust 后端性能优化和 caching 策略的参考。

讨论亮点

没有实质性的技术讨论。PR 获得 njhill 的 approve。机器人评论提示合并冲突和 pre-commit 问题,但已通过 merge 解决。

实现拆解

  1. 引入 SchedulerStatsRecorder:在 rust/src/engine-core-client/src/metrics.rs 中新增 SchedulerStatsRecorder 结构体,包含 BTreeMap<u32, SchedulerStatsHandles>。在构造时遍历所有已连接的 engine,调用 resolve_scheduler_stats_handles 预解析每个 engine 的所有固定标签 metrics handles(包括 gauge、counter、histogram 等),并存入 SchedulerStatsHandles。提供 record 方法根据 engine 索引查询 handles 并记录 stats,替换原有的 record_scheduler_stats 函数。

  2. 重构 RequestMetricsTracker:在 rust/src/llm/src/request_metrics.rs 中新增 RequestMetricHandles 结构体,缓存单个请求(模型 + engine 索引)的所有固定标签 metrics handles,包括 generation_tokens、prompt_tokens、各种时间直方图等。RequestMetricsTracker::new 现在接受额外的 engine_index 参数,并在初始化时调用 resolve_request_metric_handles 解析 handles。observe_output 等方法不再传递 engine_index,直接使用缓存的 handles 进行指标记录。

  3. 传递 engine_index:在 rust/src/engine-core-client/src/client/stream.rs 中为 EngineCoreOutputStream 添加 engine_index 字段,并在创建时设置。在 rust/src/llm/src/lib.rs 中,创建 RequestMetricsTracker 时从 GenerateOutputStream 获取 engine_index,修改了调用顺序以确保 stream 先创建。

  4. 简化输出流:在 rust/src/llm/src/output.rs 中,GenerateOutputStreampoll_next 调用 request_metrics.observe_output 时不再传递 engine_index,因为已缓存在内部 handles 中。

  5. 测试适配rust/src/llm/tests/generate.rs 调整相关调用以匹配新的构造函数,测试覆盖无功能变化。

文件 模块 状态 重要度
rust/src/llm/src/request_metrics.rs 请求指标 modified 8.2
rust/src/engine-core-client/src/metrics.rs 引擎客户端指标 modified 8.1
rust/src/engine-core-client/src/client/stream.rs 引擎客户端 modified 5.88
rust/src/engine-core-client/src/client/imp.rs 引擎客户端 modified 5.8
rust/src/llm/src/lib.rs 请求指标 modified 5.47
rust/src/llm/src/output.rs 请求指标 modified 4.67
rust/src/metrics/src/lib.rs 指标定义 modified 4.27
rust/src/engine-core-client/src/client.rs 引擎客户端 modified 3.91
rust/src/llm/tests/generate.rs 测试 modified 3.71

关键符号

resolve_request_metric_handles resolve_scheduler_stats_handles record_scheduler_stats_with_handles RequestMetricsTracker::new RequestMetricsTracker::observe_output SchedulerStatsRecorder::new SchedulerStatsRecorder::record

关键源码片段

rust/src/llm/src/request_metrics.rs core-logic

核心文件:引入 RequestMetricHandles 结构体缓存所有请求级 metrics handles,重构 RequestMetricsTracker 使其在构造时预解析 handles,所有方法直接使用缓存 handles 而非重复 get_or_create。

/// 预解析并缓存请求级 metrics handles,避免每次记录时重复 get_or_create。
#[derive(Clone)]
struct RequestMetricHandles {
    labels: EngineLabels,    // 请求级别计数器
    num_preemptions: U64Counter,
    prompt_tokens: U64Counter,
    prompt_tokens_local_compute: U64Counter,
    prompt_tokens_local_cache_hit: U64Counter,
    prompt_tokens_external_kv_transfer: U64Counter,
    prompt_tokens_cached: U64Counter,
    generation_tokens: U64Counter,    // 请求生命周期指标(直方图等)
    request_success: Family<FinishedReasonLabels, U64Counter>,
    request_prompt_tokens: HistogramMetric,
    request_generation_tokens: HistogramMetric,
    request_max_num_generation_tokens: HistogramMetric,
    request_params_max_tokens: HistogramMetric,
    request_params_n: HistogramMetric,
    request_prefill_kv_computed_tokens: HistogramMetric,
    time_to_first_token_seconds: HistogramMetric,
    inter_token_latency_seconds: HistogramMetric,
    e2e_request_latency_seconds: HistogramMetric,
    request_queue_time_seconds: HistogramMetric,
    request_prefill_time_seconds: HistogramMetric,
    request_decode_time_seconds: HistogramMetric,
    request_inference_time_seconds: HistogramMetric,
    request_time_per_output_token_seconds: HistogramMetric,
}impl RequestMetricsTracker {
    pub(crate) fn new(
        model_name: String,
        engine_index: u32, // 新增 engine_index 参数
        arrival_time: f64,
        prompt_len: u32,
        max_tokens_param: Option<u32>,
        n_param: u32,
    ) -> Self {
        Self {
            // 在构造时预解析所有 handles
            handles: resolve_request_metric_handles(&model_name, engine_index),
            arrival_time,
            prompt_len,
            max_tokens_param,
            n_param,
            is_prefilling: true,
            queued_ts: 0.0,
            scheduled_ts: 0.0,
            first_token_ts: 0.0,
            last_token_ts: 0.0,
            first_token_latency: 0.0,
            num_generation_tokens: 0,
            latest_num_cached_tokens: 0,
        }
    }    pub(crate) fn observe_output(
        &mut self,
        batch_timestamp: f64,
        received_at: f64,
        output: &EngineCoreOutput,
    ) {
        // 直接使用缓存的 handles,无需每次 get_or_create
        self.handles.generation_tokens.inc_by(output.new_token_ids.len() as u64);
        // ... 简化,不再传递 engine_index
    }
}
rust/src/engine-core-client/src/metrics.rs core-logic

核心文件:引入 SchedulerStatsRecorder 结构体,预解析所有 scheduler stats 的 handles,提供 record 方法替代原有函数。

/// 缓存 scheduler stats 的 metric handles,按 engine 索引存储。
pub(crate) struct SchedulerStatsRecorder {
    engines: BTreeMap<u32, SchedulerStatsHandles>,
}/// 单个 engine 的缓存的 metric handles。
struct SchedulerStatsHandles {
    labels: EngineLabels,
    scheduler_running: U64Gauge,
    scheduler_waiting: U64Gauge,
    scheduler_waiting_capacity: U64Gauge,
    scheduler_waiting_deferred: U64Gauge,
    kv_cache_usage: F64Gauge,
    prefix_cache_queries: U64Counter,
    prefix_cache_hits: U64Counter,
    external_prefix_cache_queries: U64Counter,
    external_prefix_cache_hits: U64Counter,
    spec_decode_num_drafts: U64Counter,
    spec_decode_num_draft_tokens: U64Counter,
    spec_decode_num_accepted_tokens: U64Counter,
    spec_decode_num_accepted_tokens_per_pos: Family<EnginePositionLabels, U64Counter>,
    estimated_flops_per_gpu: U64Counter,
    estimated_read_bytes_per_gpu: U64Counter,
    estimated_write_bytes_per_gpu: U64Counter,
    kv_block_lifetime_seconds: HistogramMetric,
    kv_block_idle_before_evict_seconds: HistogramMetric,
    kv_block_reuse_gap_seconds: HistogramMetric,
    log_stats: SchedulerLogStatsAccumulator,
}impl SchedulerStatsRecorder {
    /// 在构造时遍历所有 engine,预解析 handles。
    pub(crate) fn new(
        metrics: &SchedulerMetrics,
        model_name: &str,
        engines: &[ConnectedEngine],
    ) -> Self {
        let engines = engines
            .iter()
            .filter_map(|engine| {
                let engine = engine.engine_id.engine_index()?;
                Some((
                    engine,
                    resolve_scheduler_stats_handles(metrics, model_name, engine),
                ))
            })
            .collect();
        Self { engines }
    }    /// 记录 scheduler stats,直接使用缓存 handles。
    pub(crate) fn record(&self, engine_index: u32, stats: &SchedulerStats) {
        if let Some(handles) = self.engines.get(&engine_index) {
            record_scheduler_stats_with_handles(handles, stats);
        }
    }
}

评论区精华

合并冲突与 pre-commit other

mergify[bot] 提示合并冲突和 pre-commit 失败,要求 rebase。

结论:已通过合并 main 分支解决。 · 已解决

风险与影响

主要风险在于 handles 缓存可能导致与底层 metrics 注册表不同步,但 handles 是通过 get_or_create_owned 获取的,持有强引用,因此安全。若 engine 集合在运行时动态变化(如 DP 扩展),SchedulerStatsRecorder 需要重建,当前设计仅在 ClientInner::new 时初始化,后续变更可能不生效。需要关注 Family 类型指标的 child lookup 性能(如 spec_decode_num_accepted_tokens_per_pos 仍每次解析 position 标签)。改动涉及多个文件,需确保各接口调用匹配正确。

对用户透明,内部性能提升。减少每次 metrics 记录时的哈希查找和锁竞争。主要影响 Rust frontend 的 metrics 路径。改动覆盖面广但逻辑单纯,回归风险可控。测试文件和 CI 通过。

核心缓存路径变更 缺少基准测试验证 动态 engine 扩容不支持 动态标签仍存在 child lookup 开销

关联 Issue

未识别关联 Issue

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

完整报告

参与讨论