执行摘要
- 一句话:为SimpleCPUOffloadConnector添加reset_cache()
- 推荐动作:建议PR审核者重点验证与dedup修复(PR#41289)的配合,确认
_in_flight_store_gpu_blocks在reset中正确清理。对于KV offload领域的开发者,本PR的'abandoned transfer'设计模式值得学习,用于在异步系统中安全重置状态。一般读者可跳过实现细节,只关注决策过程。
功能与动机
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.'
实现拆解
-
引入abandoned状态字典:在SimpleCPUOffloadScheduler.__init__()中添加_abandoned_store_event_to_blocks和_abandoned_reqs_to_load字典,用于在reset()时将未完成的store和load事件移入,保证引用不会被立即释放,直到对应传输完成。
-
实现reset()方法:在manager.py中新增reset()方法。该方法将当前所有_store_event_to_blocks和_reqs_to_load移入abandoned字典,清除调度器队列、CPU前缀缓存,重置光标(lazy模式)。如果abandoned字典非空则返回False,所有pending传输完成后返回True。要求_gpu_block_pool必须已绑定,否则断言失败。
-
调整传输完成处理:修改_process_store_event(),先尝试从_store_event_to_blocks弹出transfer;若为None则从abandoned字典弹出;若仍为None则忽略(陈旧事件)。之后调用新增的_release_transfer_refs()释放CPU/GPU块引用,重置CPU块的哈希并将其归还到空闲池。
-
更新连接器接口:在SimpleCPUOffloadConnector.reset_cache()中将raise NotImplementedError替换为return self.scheduler_manager.reset(),并添加注释说明worker侧不主动接触,in-flight传输会自然完成,陈旧完成事件会被忽略。
-
添加单元测试:在test_scheduler.py中新增三个测试函数,分别验证eager模式、lazy模式和pending load场景下reset()的正确性。测试中构建in-flight传输,调用reset()确认返回False,模拟完成事件后检查块引用释放和前缀缓存可重置。
关键文件:
vllm/v1/simple_kv_offload/manager.py(模块 KV卸载;类别 source;类型 core-logic;符号 _release_transfer_refs, reset): 核心实现,包含reset逻辑、abandoned状态追踪和引用释放
tests/v1/simple_kv_offload/test_scheduler.py(模块 测试;类别 test;类型 test-coverage;符号 test_reset_pending_eager_stores, test_reset_pending_lazy_stores, test_reset_pending_loads): 新增三个测试用例,验证reset在各种模式下的正确性
vllm/distributed/kv_transfer/kv_connector/v1/simple_cpu_offload_connector.py(模块 连接器;类别 source;类型 core-logic): 将reset_cache从NotImplementedError改为委托实现
关键符号:reset, _release_transfer_refs, _process_store_event, has_pending_stores, reset_cache
关键源码片段
vllm/v1/simple_kv_offload/manager.py
核心实现,包含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
新增三个测试用例,验证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
将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
评论区精华
-
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不应在调度器未初始化时调用。作者随后改为断言。
- 确保worker侧异步传输完成后再释放资源 (design): 采用abandoned方案,等待传输自然完成,不主动同步worker。
- 使用set避免reset中重复free导致引用计数错误 (correctness): 待合并dedup修复,作者可能需要在reset中补充清理_in_flight_store_gpu_blocks。
- 将gpu_pool is None检查改为断言 (design): 改为断言,更早暴露状态错误。
- worker是否需要同步reset (design): worker侧不处理,依赖自然完成+abandoned机制。
风险与影响
- 风险:
- 依赖外部修复:reset的正确性依赖于即将合并的dedup修复(PR#41289),若未合并或修复不完整,可能导致引用计数异常。
- 异步竞争:虽然abandoned机制防止了传输中释放,但未同步worker,如果worker在reset后立即发送完成事件,abandoned字典可能已被清除,但事件被正确处理(忽略)。但若多个reset连续调用,可能造成状态混乱。
- 测试覆盖不足:测试仅覆盖单次reset,缺少多级reset嵌套或并发场景。
- 资源泄漏:如果传输永远不完成(异常情况),abandoned条目会永久存在,可能导致内存泄漏。但生产环境中传输超时会保障有限等待。
- 影响:影响范围限定在使用SimpleCPUOffloadConnector的模型推理路径。用户现在可以调用reset_cache()安全重置CPU前缀缓存,等待所有in-flight DMA传输完成后才真正释放GPU/CPU块。对其他连接器无影响,API兼容性未破坏(之前是NotImplementedError,现在是正常返回)。
- 风险标记:依赖dedup修复, 异步传输竞争, 测试覆盖有限
关联脉络
- PR #41289 Fix dedup of GPU blocks in eager offloading: ivanium在review中指出该PR修复了跨步重复bug,要求在reset中清理_in_flight_store_gpu_blocks以确保正确性。
参与讨论