Prhub

#45823 [Fix][KV offload] Defer `on_request_finished` until in-flight transfers drain

原始 PR 作者 ronensc 合并时间 2026-06-18 12:05 文件变更 4 提交数 6 评论 5 代码增减 +174 / -7

执行摘要

[KV offload] 延迟 on_request_finished 直至在途传输完成

根据 PR body,当前 submit_store() 可能在 SecondaryTierManageron_request_finished() 之后被调用:连接器在 GPU→primary 存储仍在传输中时就急切地调用 on_request_finished(),后续传输完成才驱动 complete_store()submit_store() 级联到 secondary tier。此 PR 修复该时序问题。关联 roadmap issue #33689。

此 PR 值得精读,特别是对 KV 卸载模块的正确性有直接影响。设计上采用的延迟回调模式(defer until drain)是一种常见的异步资源释放策略,值得参考。审核过程对 reset_cache 安全顺序的讨论也体现了小心处理重置时机的重要性。

讨论亮点
  • 测试函数命名:作者最初用 _run() 简化测试,审核者 orozery 提议改用公共 run() 方法,因为 async/sync 已无区别。作者合并 main 后采纳。
  • reset_cache 调用顺序orozery 建议在 manager.reset_cache() 之前先处理延迟的 on_request_finished,以避免清理后丢失状态。作者按建议修改。

实现拆解

  1. 修改 update_connector_outputvllm/distributed/.../scheduler.py):当作业完成且从 transfer_jobs 移除后,若请求已结束且 transfer_jobs 为空,则调用 self.manager.on_request_finished(req_status.req_context),并删除 _req_status 条目。
  2. 修改 request_finished:当被调度器调用时,若请求无在途作业,则立即调用 on_request_finished 并清理;若有在途作业,则延迟到 update_connector_output 处理,同时将非滑动窗口块注册到 _block_id_to_pending_jobs 防止块重用。
  3. 修改 reset_cache:在调用 self.manager.reset_cache() 之前,遍历 _req_status,对已完成且有延迟调用的请求先调用 on_request_finished 并删除其状态,防止因作业被丢弃而泄漏。
  4. 更新基类文档:在 vllm/v1/kv_offload/base.pytiering/base.py 中明确 on_request_finished 的调用时机(请求无在途传输作业时触发,但已提交的异步传输可能仍在进行)。
  5. 添加测试:新增 test_no_offload_call_after_on_request_finishedtest_reset_cache_finalizes_finished_request_with_pending_store,验证延迟语义和 reset_cache 的更正。
文件 模块 状态 重要度
vllm/distributed/kv_transfer/kv_connector/v1/offloading/scheduler.py 卸载调度器 modified 6.81
tests/v1/kv_connector/unit/offloading_connector/test_scheduler.py 卸载调度器 modified 6.76
vllm/v1/kv_offload/tiering/base.py 分层卸载 modified 4.43
vllm/v1/kv_offload/base.py 卸载基类 modified 4.4

关键符号

OffloadingConnectorScheduler.update_connector_output OffloadingConnectorScheduler.request_finished OffloadingConnectorScheduler.reset_cache test_no_offload_call_after_on_request_finished test_reset_cache_finalizes_finished_request_with_pending_store

关键源码片段

vllm/distributed/kv_transfer/kv_connector/v1/offloading/scheduler.py core-logic

核心逻辑修改:实现 on_request_finished 延迟调用和 reset_cache 泄漏修复。

# scheduler.py - update_connector_output 中的延迟调用
    if not req_status.transfer_jobs and req_status.req.is_finished():
        self.manager.on_request_finished(req_status.req_context)
        del self._req_status[job_status.req_id]# scheduler.py - request_finished 延迟逻辑
    if not req_status.transfer_jobs:
        self.manager.on_request_finished(req_status.req_context)
        del self._req_status[request.request_id]
        return False, None
    # 注册 pending blocks...# scheduler.py - reset_cache 中完成延迟请求
    for req_id, status in list(self._req_status.items()):
        if status.req.is_finished():
            self.manager.on_request_finished(status.req_context)
            del self._req_status[req_id]
    self.manager.reset_cache()

评论区精华

测试使用 _run() vs run() 测试

作者在 review 中询问是否可以用 _run() 代替 run(),因为该测试只关心调用顺序,不需要块计算验证。审核者 orozery 建议改用 run() 因为 async/sync 不再有区别,且 prefer public API。

结论:最终使用了 run()。 · 已解决

reset_cache 中 on_request_finished 的调用顺序 正确性

审核者 orozery 建议在 manager.reset_cache() 之前触发延迟的 on_request_finished,以避免重置后丢失状态。作者采纳并修改代码。

结论:按建议修改:在 reset_cache 中先处理已完成请求的延迟调用,再调用 manager.reset_cache()。 · 已解决

风险与影响

  • 时序敏感:如果 update_connector_output 中判断 not req_status.transfer_jobs and req_status.req.is_finished() 的条件有误,可能导致 on_request_finished 被跳过或重复。但测试覆盖了正常路径。
  • reset_cache 变更:在 manager.reset_cache() 前调用 on_request_finished 可能触发额外的副作用,但该调用本身只应释放簿记,不应有 I/O 操作(按文档)。风险较低。
  • 影响范围:变更局限在 KV 卸载连接器的 scheduler 和基类 docstring,不影响其他子系统。
  • 用户:修复了次级 tier 存储在请求结束后仍意外调用的 bug,提高卸载可靠性和数据一致性。
  • 系统reset_cache 不再泄漏请求状态,减少了内存使用(但影响微小)。
  • 团队:代码基语义更清晰,基类文档明确了 on_request_finished 的契约。
时序依赖 reset_cache 泄漏修复 核心调度路径变更

关联 Issue

#33689 [RFC]: KV Offloading Roadmap

完整报告

参与讨论