Prhub

#48798 Add tiering offloading metrics

原始 PR 作者 Srinivasoo7 合并时间 2026-08-10 10:23 文件变更 16 提交数 13 评论 61 代码增减 +1031 / -237

执行摘要

为 KV 多级卸载新增分层指标,覆盖传输、查询、失败与活跃度

PR body 明确该变更 'Covers the scope of TieringOffloadingSpec-Level Metrics in #44008',目标是为多级 KV 卸载提供规范级可观测性:一是新增带 per-tier 标签的指标定义,二是让 JobResult 携带传输大小/时间,三是从 TieringOffloadingManager 发出读写、失败与 block 查询/命中计数。此前该子系统只有 LOOKUP_SYNC_DELAY 与 LOOKUP_ASYNC_DELAY 两个直方图,缺少传输量、失败率与 tier 占用情况,运维无法判断各 tier 的健康度与瓶颈。

值得精读。重点学习:可观测性数据如何跨异步线程传递(JobResult.transfer_time)、追踪器与 orchestrator 的职责划分、per-block 异步延迟的语义定义,以及 NIXL telemetry 的接入方式。后续关注 #44008 的剩余指标项与 p2p telemetry follow-up。

讨论亮点

orozery 是主要审核者,围绕设计做了多轮交锋并最终 APPROVED。核心讨论:

  • 抽取追踪器:orozery 认为指标代码 (~130 行) 有凝聚力,建议挪到 tiering/metrics.py,使 manager.py 保持纯编排、tracker 可独立测试,落地为 TieringMetricsTracker
  • lookup 延迟语义:orozery 指出 per-request 异步延迟会把 promotion 时间也算进去,建议按 per-block/per-tier 记录首次 unresolved 起点、解析时再观测,最终实现 observed_lookups 存 start time 并配直方图。
  • 计数去重:orozery 要求 queries/hit 只计一次 per [req][block][tier],且分配后停止,最终通过 on_request_allocatedobserved_lookups=None 实现。
  • 传输耗时度量:orozery 要求 read/write time 只算纯 I/O 时间,启发 fs 线程池在 task() 前后计时;obj/p2p 改用 agent.get_xfer_telemetry(handle),p2p 因需要额外 plumbing 留作 @liranschour 的 follow-up。
  • 数据契约简化:orozery 建议移除 JobResult.transfer_size,由 manager 按 key 数与 primary block size 计算。

实现拆解

1. 指标注册与语义重构(spec.py)

TieringOffloadingSpec.build_metric_definitions() 新增 12 个指标,除 PROMOTION_ALLOCATION_FAILURES 外全部带 ("tier",) 标签;同时把既有 LOOKUP_SYNC_DELAY/LOOKUP_ASYNC_DELAY 直方图的语义从请求级累计改为 per-block、per-tier 的阻塞/异步等待,并同步更新 documentation。

2. 传输契约扩展(base.py + fs/obj/example/p2p 四类 tier)

JobMetadata 更名为 TransferJobJobResult 新增 transfer_time 字段;fs 的 DualQueueThreadPool._worker 在 task() 前后计时,get_finished 返回累计时长;obj 改用 agent.get_xfer_telemetry(handle) 获取真实传输时长;example tier 填 0;p2p 暂用本地时间差并留待后续接入 telemetry。

3. 追踪器落地(metrics.py,新增)

TieringMetricsTracker 集中维护 _RequestMetricsState(lookup 去重与异步起点)和 _TierState(活跃 job/读写块数增量计数),提供 on_lookupon_job_registeredon_job_finishedtake_statsassert_idle 等接口;lookup 计数按 (request, block, tier) 去重,请求分配后停止上报。

4. Manager 集成与重构(manager.py)

引入 JobMetadata(transfer_job, tier_idx) NamedTuple 统一 job 跟踪,_register_job/_pop_job 成为唯一入口并同步 metrics;_SecondaryTierFacingParent_pending_load_submissionsrequest_level_tiers 全部从 tier 对象改为 tier_idx,消除重复映射;get_stats() 聚合 tracker 与各 secondary tier 的统计。

5. 测试配套

新增 tests/v1/kv_offload/tiering/test_metrics.py(6 个用例覆盖 lookup、去重、job 完成、部分成功、gauge、分配失败);改造 tests/v1/kv_offload/tiering/test_tiering_offloading.py 与 tests/v1/kv_connector/unit/offloading_connector/test_metrics.py 验证 spec 注册与 manager 聚合;各 tier 测试同步适配 TransferJob 接口。

文件 模块 状态 重要度
vllm/v1/kv_offload/tiering/metrics.py 指标层 added 8.98
vllm/v1/kv_offload/tiering/manager.py 多级卸载 modified 8.69
vllm/v1/kv_offload/tiering/base.py 数据契约 modified 7.11
vllm/v1/kv_offload/tiering/spec.py 指标注册 modified 6.15
vllm/v1/kv_offload/tiering/fs/thread_pool.py 线程池 modified 6.84
tests/v1/kv_offload/tiering/test_metrics.py 指标测试 added 7.18

关键符号

TieringMetricsTracker.__init__ TieringMetricsTracker.on_lookup TieringMetricsTracker.on_job_registered TieringMetricsTracker.on_job_finished TieringMetricsTracker.take_stats TieringOffloadingManager._register_job TieringOffloadingManager._pop_job TieringOffloadingManager.get_stats DualQueueThreadPool._worker JobState.task_done TieringOffloadingSpec.build_metric_definitions

关键源码片段

vllm/v1/kv_offload/tiering/manager.py core-logic

指标集成与核心重构发生地:统一 job 跟踪、接入 tracker、parent 包装改用 tier_idx,直接影响整个卸载协调流程。

# vllm/v1/kv_offload/tiering/manager.py —— job 注册 / 弹出与 stats 聚合(整理后)class JobMetadata(NamedTuple):
    # transfer_job 描述一次异步传输;tier_idx 记录该 job 归属的 secondary tier
    transfer_job: TransferJob
    tier_idx: intclass TieringOffloadingManager(OffloadingManager):
    def _register_job(self, transfer_job: TransferJob, tier_idx: int) -> None:
        # 统一入口:登记 job 元数据的同时更新 tracker 的活跃计数与 gauge
        job_metadata = JobMetadata(transfer_job, tier_idx)
        self._jobs[transfer_job.job_id] = job_metadata
        self._metrics.on_job_registered(job_metadata)
​
    def _pop_job(self, job_id: JobId) -> JobMetadata | None:
        # 统一出口:job 完成时移除登记并让 tracker 同步递减活跃状态
        return self._jobs.pop(job_id, None)
​
    def get_stats(self) -> OffloadingConnectorStats | None:
        # 聚合 tracker 的观察结果与各 secondary tier 自行上报的统计;
        # 返回后由上层收集为 Prometheus 指标
        stats = self._metrics.take_stats()
        for tier in self.secondary_tiers:
            tier_stats = tier.get_stats()
            if tier_stats is not None:
                if stats is None:
                    stats = tier_stats
                else:
                    stats.aggregate(tier_stats)
        return stats
vllm/v1/kv_offload/tiering/base.py data-contract

定义传输数据契约:JobMetadata 更名为 TransferJob,JobResult 增加 transfer_time,新增全部指标常量,是各 tier 实现的公共接口。

# vllm/v1/kv_offload/tiering/base.py —— 异步传输数据契约(整理后)@dataclass
class TransferJob:
    # Metadata for an in-flight async transfer job.
    job_id: JobId
    keys: Collection[OffloadKey]
    block_ids: np.ndarray
    is_promotion: bool # True:secondary → primary(promotion);False:primary → secondary(cascade)
    req_context: ReqContext
​
​
@dataclass
class JobResult:
    # Result of an async transfer job.
    job_id: JobId
    success: bool
    # 仅 promotion 部分失败时使用,标识成功加载的 keys;None 表示全部一致
    successful_keys: Collection[OffloadKey] | None = None
    # 纯 I/O 传输耗时(秒):fs 来自线程池 task() 前后计时,obj 来自 NIXL telemetry
    transfer_time: float | None = None

评论区精华

将指标逻辑从 manager.py 抽离为 TieringMetricsTracker 设计

orozery 认为 manager.py 中指标代码 (~130 行 ) 有凝聚力,建议抽到 tiering/metrics.py 使 manager 保持纯编排、tracker 可独立测试。

结论:已新建 metrics.py 并移入全部指标状态与方法,落地为 TieringMetricsTracker。 · 已解决

lookup 异步延迟度量改为 per-block per-tier 设计

orozery 指出 per-request 异步延迟会把 promotion 时间也算进去,语义不清晰;建议参照 offloading connector 的 _maybe_observe_lookup_async_delay,按 request/tier/block 记录首次 unresolved 起点,resolved 时观测。

结论:已实现 observed_lookups 存 start time,sync/async 延迟均 per-block per-tier 上报。 · 已解决

每 [req][block][tier] 的查询 / 命中计数去重 正确性

orozery 要求 queries/hit 只计一次,且请求分配后停止上报,避免重复统计与分配后噪音。

结论:on_lookup 中 setdefault 标记已观测,on_request_allocated 置 observed_lookups 为 None;新增专项测试覆盖。 · 已解决

fs tier 的 read/write time 只统计 IO 时间 性能

orozery 指出传输时间不应包含排队、等待收集的时间,应在 DualQueueThreadPool._worker 的 task() 前后计时。

结论:已改在任务线程内计时,get_finished 返回累计 transfer_time。 · 已解决

obj/p2p tier 使用 NIXL telemetry 获取传输时长 设计

orozery 建议用 agent.get_xfer_telemetry(handle) 而非本地时间戳,p2p 需要 DataTransport 额外 plumbing。

结论:obj 已接入 telemetry;p2p 保留为 @liranschour 的 follow-up。 · 已解决

JobResult.transfer_size 改由 manager 计算 设计

orozery 建议移除 JobResult.transfer_size,由 manager 按 key 数与 CPU page size 计算,避免各 tier 重复上报。

结论:已移除,transfer_size = completed_key_count * primary_block_size。 · 已解决

create_store_job 的 tier_idx 可选参数 style

orozery 要求移除 tier_idx 的 None 默认值,调用处 assert,避免可选参数掩盖错误。

结论:tier_idx 改为必传。 · 已解决

风险与影响

  • 计数状态一致性:_TierState 的增减必须在 job 注册/完成时严格配对,assert_idle 与内部断言是主要防线;若未来 tier 实现绕过 _register_job/_pop_job 会静默破坏 gauge。
  • 接口兼容性:JobMetadataTransferJob 更名与 JobResult.transfer_time 新增影响所有 SecondaryTierManager 子类(fs/obj/example/p2p),第三方自定义 tier 需要同步适配。
  • 热路径开销:每次 lookup 增加 time.monotonic() 与字典操作,但仅在请求分配前执行且每 block 只计入一次;量级可忽略,仍建议在大规模 prefix-cache 场景跑基准确认。
  • 精度边界:fs 的 transfer_time 是多任务累计值;p2p 仍是本地时间差而非真实传输时长,跨机场景偏差可能较大。
  • 测试盲区:指标注册与 tracker 行为覆盖充分,但未做 Prometheus scraping 的端到端断言。
  • 用户/运维:新增 vllm:kv_offload_tiering_* 系列指标,可直接观测各 tier 的读写量、耗时、失败率与 primary 占用,便于定位卸载瓶颈。
  • 系统:lookup 与 job 生命周期路径新增轻量计数/计时,对推理主路径无正确性影响。
  • 团队:5 个 tier 实现需同步接口,metrics.py 成为后续指标接入的范式(tracker 模式 + telemetry 取时长)。
  • 该变更限定在 v1 kv_offload 子系统内部,不触碰调度、注意力等核心路径。
接口重构波及 4 个 tier 实现 lookup 热路径新增计时与字典操作 计数状态依赖断言维护 P2P 时长暂用近似值 指标链路缺端到端验证

关联 Issue

未识别关联 Issue

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

完整报告

参与讨论