Prhub

#39726 [SimpleCPUOffloadConnector]: Add support for reset_cache()

原始 PR 作者 jonathanc-n 合并时间 2026-06-18 10:47 文件变更 3 提交数 10 评论 16 代码增减 +251 / -9

执行摘要

为 SimpleCPUOffloadConnector 添加 reset_cache()

PR描述指出,Scheduler.reset_prefix_cache()在in-flight DMA传输时失败,因为ref_cnt>0的块仍被使用。本PR旨在添加reset_cache()支持,通过设计让所有改动局限在SimpleCPUOffloadConnector内,避免与worker同步,如pr body所述:'These changes purposely are done this way to allow zero changes outside of the SimpleCPUOffloadConnector and not need to sync with the worker.'

建议PR审核者重点验证与dedup修复(PR#41289)的配合,确认_in_flight_store_gpu_blocks在reset中正确清理。对于KV offload领域的开发者,本PR的'abandoned transfer'设计模式值得学习,用于在异步系统中安全重置状态。一般读者可跳过实现细节,只关注决策过程。

讨论亮点
  • aoshen02 对同步安全性的质疑:'仅清除调度器侧状态不足,worker侧异步传输可能仍在运行。' jonathanc-n 通过引入abandoned字典保证引用不释放,reset在传输完成前返回False。方案被接受。

  • gemini-code-assist[bot] 关于重复free的建议:使用set收集block ID防止重复free。jonathanc-n认为当前逻辑已保证去重,ivanium指出存在跨步重复bug,将在PR#41289中修复,并建议在reset中清理_in_flight_store_gpu_blocks。该建议被部分采纳,等待dedup修复合并。

  • NickLucche 对保护条件的建议:将gpu_pool is None的静默返回改为断言,因为reset不应在调度器未初始化时调用。作者随后改为断言。

实现拆解

  1. 引入abandoned状态字典:在SimpleCPUOffloadScheduler.__init__()中添加_abandoned_store_event_to_blocks_abandoned_reqs_to_load字典,用于在reset()时将未完成的store和load事件移入,保证引用不会被立即释放,直到对应传输完成。

  2. 实现reset()方法:在manager.py中新增reset()方法。该方法将当前所有_store_event_to_blocks_reqs_to_load移入abandoned字典,清除调度器队列、CPU前缀缓存,重置光标(lazy模式)。如果abandoned字典非空则返回False,所有pending传输完成后返回True。要求_gpu_block_pool必须已绑定,否则断言失败。

  3. 调整传输完成处理:修改_process_store_event(),先尝试从_store_event_to_blocks弹出transfer;若为None则从abandoned字典弹出;若仍为None则忽略(陈旧事件)。之后调用新增的_release_transfer_refs()释放CPU/GPU块引用,重置CPU块的哈希并将其归还到空闲池。

  4. 更新连接器接口:在SimpleCPUOffloadConnector.reset_cache()中将raise NotImplementedError替换为return self.scheduler_manager.reset(),并添加注释说明worker侧不主动接触,in-flight传输会自然完成,陈旧完成事件会被忽略。

  5. 添加单元测试:在test_scheduler.py中新增三个测试函数,分别验证eager模式、lazy模式和pending load场景下reset()的正确性。测试中构建in-flight传输,调用reset()确认返回False,模拟完成事件后检查块引用释放和前缀缓存可重置。

文件 模块 状态 重要度
vllm/v1/simple_kv_offload/manager.py KV 卸载 modified 7.78
tests/v1/simple_kv_offload/test_scheduler.py 测试 modified 7.39
vllm/distributed/kv_transfer/kv_connector/v1/simple_cpu_offload_connector.py 连接器 modified 5.97

关键符号

reset _release_transfer_refs _process_store_event has_pending_stores reset_cache

关键源码片段

vllm/v1/simple_kv_offload/manager.py core-logic

核心实现,包含 reset 逻辑、abandoned 状态追踪和引用释放

    def reset(self) -> bool:
        """Abort all pending transfers, release block refs, reset CPU cache.
        Returns False if there are in-flight DMA transfers that must complete first.
        """
        gpu_pool = self._gpu_block_pool
        assert gpu_pool is not None, "reset() called before GPU pool bound"
​
        # 1. Move all pending store events to abandoned set
        # (their refs remain pinned until worker confirms completion)
        self._abandoned_store_event_to_blocks.update(
            self._store_event_to_blocks
        )
        self._store_event_to_blocks.clear()
​
        # 2. Move all pending loads to abandoned set
        self._abandoned_reqs_to_load.update(self._reqs_to_load)
        self._reqs_to_load.clear()
        self._load_event_to_reqs.clear()
​
        # 3. Clear scheduler-side queues and reset CPU prefix cache
        self._pending_cpu_hits.clear()
        self.cpu_block_pool.reset_prefix_cache()
​
        # 4. In lazy mode, reset the cursor
        if self._lazy_mode:
            self._cursor = None
​
        # 5. Return True iff no in-flight transfers remain
        return not (self._abandoned_store_event_to_blocks
                    or self._abandoned_reqs_to_load)
​
    def _release_transfer_refs(self, transfer: TransferMeta) -> None:
        """Release GPU/CPU block refs after a transfer (store or load) completes.
        CPU block hashes are cleared so the data won't be treated as valid cache.
        """
        cpu_blocks = [self.cpu_block_pool.blocks[bid] for bid in transfer.cpu_block_ids]
        for cb in cpu_blocks:
            cb.reset_hash()
        self.cpu_block_pool.free_blocks(cpu_blocks)
​
        assert self._gpu_block_pool is not None
        self._gpu_block_pool.free_blocks(
            self._gpu_block_pool.blocks[bid] for bid in transfer.gpu_block_ids
        )
tests/v1/simple_kv_offload/test_scheduler.py test-coverage

新增三个测试用例,验证 reset 在各种模式下的正确性

    # ------------------------------------------------------------------
    # Test 12: Reset with pending eager stores waits for completion
    # ------------------------------------------------------------------
    def test_reset_pending_eager_stores() -> None:
        """Eager mode: reset() abandons in-flight stores until they complete."""
        fix = make_scheduler(num_cpu_blocks=8, num_gpu_blocks=16, lazy=False)
        sched = fix.scheduler
        gpu_pool = fix.gpu_block_pool
​
        num_blocks = 2
        req = make_request(num_blocks=num_blocks)
​
        kv_blocks = _alloc_and_register(fix, req, num_blocks)
        sched.update_state_after_alloc(req, kv_blocks, num_external_tokens=0)
        block_ids = kv_blocks.get_block_ids()
        sched_out = make_scheduler_output(
            {req.request_id: num_blocks * BLOCK_SIZE},
            new_reqs={req.request_id: block_ids},
        )
​
        meta = sched.build_connector_meta(sched_out)
        assert meta.store_event >= 0
        assert len(sched._store_event_to_blocks) > 0
​
        # GPU blocks should have elevated ref_cnt from touch()
        for bid in meta.store_gpu_blocks:
            assert gpu_pool.blocks[bid].ref_cnt > 0
​
        # Free the request's own block refs (simulates preemption)
        gpu_pool.free_blocks(gpu_pool.blocks[bid] for bid in block_ids[0])
​
        # Reset should keep DMA refs pinned until the worker reports completion.
        assert sched.reset() is False
        assert len(sched._store_event_to_blocks) == 0
        assert len(sched._abandoned_store_event_to_blocks) == 1
        assert len(sched._reqs_to_store) == 0
        assert len(sched._store_event_to_reqs) == 0
​
        num_used = gpu_pool.num_gpu_blocks - gpu_pool.get_num_free_blocks()
        assert num_used > 1
​
        simulate_store_completion(sched, meta.store_event)
        assert len(sched._abandoned_store_event_to_blocks) == 0
​
        # All GPU blocks should now be free (ref_cnt == 0) except null block
        num_used = gpu_pool.num_gpu_blocks - gpu_pool.get_num_free_blocks()
        assert num_used == 1, f"Expected only null block in use, got {num_used}"
​
        # GPU prefix cache reset should now succeed
        assert gpu_pool.reset_prefix_cache() is True
        assert sched.reset() is True
vllm/distributed/kv_transfer/kv_connector/v1/simple_cpu_offload_connector.py core-logic

将 reset_cache 从 NotImplementedError 改为委托实现

    # NOTE: Workers are not contacted. In-flight transfers drain naturally,
    # and stale completions are ignored by the guarded
    # SimpleCPUOffloadScheduler._process_store_event().
    def reset_cache(self) -> bool | None:
        if self.scheduler_manager is not None:
            return self.scheduler_manager.reset()
        return None

评论区精华

确保 worker 侧异步传输完成后再释放资源 设计

aoshen02 评论:仅清除调度器状态不够,worker 侧 CUDA 流上的 load/store 可能仍在进行。jonathanc-n 回复:通过 abandoned 字典避免释放引用,reset 在传输完成前返回 False,传输完成后自动清理。

结论:采用 abandoned 方案,等待传输自然完成,不主动同步 worker。 · 已解决

使用 set 避免 reset 中重复 free 导致引用计数错误 正确性

gemini-code-assist[bot] 建议使用 set 收集 block ID 再 free。jonathanc-n 认为当前逻辑已保证去重。ivanium 指出存在跨步重复 bug,已提交 PR#41289 修复,并建议在 reset 中清理 _in_flight_store_gpu_blocks。

结论:待合并 dedup 修复,作者可能需要在 reset 中补充清理 _in_flight_store_gpu_blocks。 · partially resolved

将 gpu_pool is None 检查改为断言 设计

NickLucche 建议使用断言而非静默返回 True,因为 reset 不应在调度器未初始化时调用。作者随后改为断言。

结论:改为断言,更早暴露状态错误。 · 已解决

worker 是否需要同步 reset 设计

aoshen02 询问为何 worker 也会调用 reset,jonathanc 解释 worker 侧 no-op,让 in-flight 传输自然完成。

结论:worker 侧不处理,依赖自然完成 +abandoned 机制。 · 已解决

风险与影响

  1. 依赖外部修复:reset的正确性依赖于即将合并的dedup修复(PR#41289),若未合并或修复不完整,可能导致引用计数异常。
  2. 异步竞争:虽然abandoned机制防止了传输中释放,但未同步worker,如果worker在reset后立即发送完成事件,abandoned字典可能已被清除,但事件被正确处理(忽略)。但若多个reset连续调用,可能造成状态混乱。
  3. 测试覆盖不足:测试仅覆盖单次reset,缺少多级reset嵌套或并发场景。
  4. 资源泄漏:如果传输永远不完成(异常情况),abandoned条目会永久存在,可能导致内存泄漏。但生产环境中传输超时会保障有限等待。

影响范围限定在使用SimpleCPUOffloadConnector的模型推理路径。用户现在可以调用reset_cache()安全重置CPU前缀缓存,等待所有in-flight DMA传输完成后才真正释放GPU/CPU块。对其他连接器无影响,API兼容性未破坏(之前是NotImplementedError,现在是正常返回)。

依赖 dedup 修复 异步传输竞争 测试覆盖有限

关联 Issue

未识别关联 Issue

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

完整报告

参与讨论