PR vllm-project/vllm#48596 ——《[Bugfix][KV Offloading] Offload last block at request finish and prevent reuse race》。 +116/-31,改 scheduler.pyworker.py 与测试三个文件。

1. 背景:store job 的正常生命周期

OffloadingConnectorScheduler 把已计算的 KV block 从 GPU offload 到 CPU/磁盘, 供后续相同前缀的请求命中复用。核心链路是跨 step 的

schedule() 阶段
  └─ _build_store_jobs(scheduler_output)
       遍历 num_scheduled_tokens 里的请求
       对"已填满"的 block 调 manager.prepare_store() 建 store job
       job 加入 _jobs 和 req_status.transfer_jobs
  └─ build_connector_meta() → 把 store job 打进 metadata 发给 worker
下一步 worker
  └─ submit_store() 真正发起 GPU→CPU 传输

关键约束:_build_store_jobs 只处理填满的 block,未满的 partial block 会被跳过 (因为 block hash 还没定,无法作为 prefix cache 的 key)。

2. Bug 原理

_build_store_jobsschedule() 时运行——此刻结束 token(EOS)还没被处理, 所以请求的最后一个 block 仍是 partial → 被跳过。

request_finished 在 EOS 被 append 之后才调用,此时最后 block 已填满。但旧逻辑里 request_finished 没有任何 store 触发逻辑,且当请求没有在途 job 时会立刻:

del self._req_status[request.request_id]   # 旧逻辑:直接删除

→ 最后一个 block 的状态被静默丢弃,永远不会建 store job。

后果链:

请求自然结束(EOS 填满最后 block)
  └─ request_finished 立即删 _req_status
       └─ 最后 block 从未 offload
            └─ 下一个相同前缀的请求 lookup 时 miss 这个 block
                 └─ 额外 prefill 计算 + 最后 block 无法 KV 复用

这不是边角情况:每个自然结束的请求都会在 finish 时刚好填满一个 block(EOS 那一格), 所以这个 miss 是”必然可达”的。

3. 修复(scheduler 侧)

(a) request_finished 改为保活 + 刷新 key

删掉 orozery 遗留的 TODO: possibly kickoff offload for last block,改为:

self.manager.on_request_finished(req_status.req_context)
req_status.update_offload_keys()   # 用最终 block hash 刷新 offload keys
# 不再 del _req_status —— 始终保活
for job_id in req_status.transfer_jobs:      # 在途 job 的 block 注册到 flush 检测
    for bid in self._jobs[job_id].non_sliding_window_block_ids or ():
        self._block_id_to_pending_jobs.setdefault(bid, set()).add(job_id)
return False, None

(b) _build_store_jobs 纳入已结束请求

for req_id in chain(
    scheduler_output.num_scheduled_tokens,
    scheduler_output.finished_req_ids or (),   # 新增:已结束请求也处理
):
    ...
    if req.is_finished():
        num_tokens_after_batch = req.num_tokens   # 最后 block 此时已填满
    else:
        num_tokens_after_batch = req.num_computed_tokens + num_scheduled_tokens
    num_offloadable_tokens = self._calc_num_offloadable_tokens(
        req_status, num_tokens_after_batch)

于是已结束请求的最后 block 走正常代码路径建 store job,无需特判。

(c) 清理时机后移

保活之后必须有新的删除时机。抽出 helper:

def _maybe_cleanup_finished_req(self, req_id, req_status):
    if req_status.req.is_finished() and not req_status.transfer_jobs:
        del self._req_status[req_id]

_build_store_jobs 的三个”不产生 job”分支里调用(无新 offload_keys / prepare_store 分配失败 / keys_to_store 为空),保证 req_status 在正确时点释放, 不泄漏。呼应 issue #47107 描述的路径:在 request_finished 建 store job、fence block、 下一次 build_connector_meta flush

4. 修复(worker 侧 fence 竞态)

保活方案引入一个新竞态。已结束请求的 store job 是 deferred 的:

step N   : request_finished 保活
step N+1 : _build_store_jobs 建 store job(进 metadata)
step N+2 : worker 才 submit_store

若 block 在 submit 之前就被 KV cache manager 回收、复用给新请求, handle_preemptions 里的 wait() 会在 jobs_to_flush 中找到该 job,但它还没被 submit_storewait() 是 no-op → 无法 fence → block 被新数据覆盖,旧数据损坏

修复:handle_preemptionswait() 之前,先把 jobs_to_flush 的 store job 从 store_jobs 弹进 _unsubmitted_store_jobs 并让既有提交循环先提交:

if kv_connector_metadata.jobs_to_flush:
    for job_id in kv_connector_metadata.jobs_to_flush:
        entry = kv_connector_metadata.store_jobs.pop(job_id, None)
        if entry is not None:
            self._unsubmitted_store_jobs.append(
                (job_id, entry.src_spec, entry.dst_spec))
for job_id, src_spec, dst_spec in self._unsubmitted_store_jobs:
    self.worker.submit_store(job_id, src_spec, dst_spec)   # 先提交
self._unsubmitted_store_jobs.clear()
# 之后 wait() 才能真正 fence 住 block

只影响 block 被复用的场景;正常 store 仍保持 deferred,不拖慢 token 生成性能—— 这是一次典型的”正确性 vs 性能”权衡:默认走快路径,只在真正会冲突时才强制同步。

5. 跨 step 时序图

step N     : 生成 EOS,填满最后 block
             └─ request_finished(): update_offload_keys() + 保活 req_status
step N+1   : _build_store_jobs() 处理 finished_req_ids
             ├─ 建 store job(最后 block 现在已满)
             └─ 若 block 已被复用 → 加入 _current_batch_jobs_to_flush
             build_connector_meta() 发给 worker
step N+2   : worker.submit_store()
  或抢占时 : handle_preemptions() 先 submit jobs_to_flush → wait() 真正 fence

6. 测试要点

7. 小结