Prhub

#51622 [Bugfix][KV Offload] Centralize shared mmap cleanup in CPU worker

原始 PR 作者 Alex-ai-future 合并时间 2026-08-11 18:22 文件变更 2 提交数 3 评论 8 代码增减 +220 / -13

执行摘要

CPU KV 卸载统一 mmap 清理所有权,修复关停挂死

PR body 指出,CPUOffloadingWorker 创建一套 mmap 支撑的 CPU tensor 并同时传给 GPU-to-CPU(store)与 CPU-to-GPU(load)两个 handler,但此前 "the region cleanup responsibility was attached to one direction, even though both directions could still have asynchronous transfers using the same backing memory",导致三类具体失败:CPU->GPU 方向仍有 in-flight 传输时 region 被提前释放;handler 在 shutdown 中抛异常会使整个 teardown 中断,另一方向和 region 都无法清理;end_event.synchronize() 因 CUDA driver fault、device lost 或从未记录的 stream 无限阻塞,engine teardown 永久挂起。目标是让 worker 成为共享 region 的唯一所有者,并保证关停流程总是能完成、失败可观测。

值得精读。该 PR 展示了异步资源生命周期的一个小而关键的所有权修复:将共享 mmap 的所有权集中到组合器(worker),方向处理器只负责自身资源;错误路径用“记录日志 + 继续执行”保证 teardown 必然完成。对从事分布式/异步资源管理的工程师有参考价值,测试手法也值得借鉴。建议跟进其配套 PR#51227(构造失败路径),确保两条生命周期路径都被覆盖。

讨论亮点

review 主要由维护者 orozery 提出 3 条意见:① 在 SingleDirectionOffloadingHandler.shutdown() 中建议同步失败时仅记录日志、清空队列并 break,不保留 sync_error 重抛;② 在 test_transfer 中明确要求保留 del 引用块,因为测试局部引用在 worker.shutdown() 期间持有 buffer exports;③ 建议 handler 关停失败后、mmap cleanup 前先强制设备同步,防止 mmap 卸载与 in-flight DMA 竞争。最终版本采纳了②和③(并用平台无关的 torch.accelerator.synchronize() 实现),①未完全采纳:最终代码保留 sync_error 并重新抛出,使其成为 worker 层感知失败、触发设备同步与日志记录的信号。

实现拆解

  1. mmap region 所有权上移SingleDirectionOffloadingHandler.__init__ 删除 mmap_region 参数与 self._mmap_region 字段(vllm/v1/kv_offload/cpu/gpu_worker.py);CPUOffloadingWorker.__init__ 在构造早期直接 self._mmap_region = mmap_region。原因是 region 被两个 handler 共享,挂在单向 handler 上正是释放时机错误之源,所有权应归属组合器 worker。
  2. 方向 handler 关停逻辑加固SingleDirectionOffloadingHandler.shutdown() 改为先看队头 self._transfers[0]end_event.synchronize() 成功后才 popleft();同步抛异常时清空队列、记录日志并保存 sync_error,避免 teardown 无限阻塞。无论成功失败,都会无条件清理 _transfer_events、stream/event/buffer pool 与 src_tensors/dst_tensors,最后重新抛出 sync_error 供上层感知。
  3. worker 统一关停编排CPUOffloadingWorker.shutdown()_store_handler_load_handler 分别用 try/except 关停,失败仅记录日志并置 handler_failed;若发生失败,先执行 torch.accelerator.synchronize() 防御 in-flight DMA 与 unpin/unmap 竞争,再统一调用一次 mmap_region.cleanup(),最后将 _mmap_regionNone
  4. 测试配套tests/v1/kv_offload/cpu/test_gpu_worker.py 新增 4 个单元测试,利用 CPUOffloadingWorker.__new__ 绕过构造直接注入 MagicMock handler,断言关停调用顺序 ["store", "load", "region"]、单方向失败后另一方向与 region 仍被清理、事件同步失败后队列与各池清空、设备同步失败时降级警告;test_transfer 移除手工 mmap_region.cleanup() 调用,改由 worker.shutdown() 统一负责,并保留 del 引用释放块。
文件 模块 状态 重要度
vllm/v1/kv_offload/cpu/gpu_worker.py KV 卸载 modified 6.66
tests/v1/kv_offload/cpu/test_gpu_worker.py KV 卸载 modified 7.06

关键符号

SingleDirectionOffloadingHandler.shutdown SingleDirectionOffloadingHandler.__init__ CPUOffloadingWorker.shutdown CPUOffloadingWorker.__init__

关键源码片段

vllm/v1/kv_offload/cpu/gpu_worker.py core-logic

核心源码改动:mmap region 所有权从单向 handler 上移到 worker,handler 只负责 drain 与传输侧资源释放,worker 统一编排两方向关停并做设备同步兜底。

class SingleDirectionOffloadingHandler:
    # ...(省略构造与传输逻辑)
​
    def shutdown(self) -> None:
        """Drain this direction and release its transfer-side resources."""
        sync_error: Exception | None = None
        while self._transfers:
            transfer = self._transfers[0]
            try:
                # 先同步队头传输的 end event,成功后才弹出该传输。
                # 若同步失败(driver fault、device lost、stream 未记录事件),
                # 直接清空队列并记录异常,避免 teardown 无限阻塞。
                transfer.end_event.synchronize()
            except Exception as e:
                logger.exception(
                    "Failed to synchronize transfer end event; "
                    "skipping %d remaining transfers",
                    len(self._transfers) - 1,
                )
                self._transfers.clear()
                sync_error = e
                break
            self._transfers.popleft()
​
        # 无论同步是否成功,都释放传输侧资源(stream、event、descriptor buffer)。
        self._transfer_events.clear()
        self._stream_pool.clear()
        self._event_pool.clear()
        self._buffer_pool.clear()
        self.src_tensors.clear()
        self.dst_tensors.clear()
        if sync_error is not None:
            raise sync_error
​
​
class CPUOffloadingWorker(OffloadingWorker):
    def shutdown(self) -> None:
        handler_failed = False
        # 两个方向的 handler 共享同一套 mmap 支撑的 CPU tensor,
        # 因此必须都尝试关停,且任一失败都不能阻止 region 清理。
        try:
            self._store_handler.shutdown()
        except Exception:
            logger.exception("Failed to shut down store offloading handler")
            handler_failed = True
​
        try:
            self._load_handler.shutdown()
        except Exception:
            logger.exception("Failed to shut down load offloading handler")
            handler_failed = True
​
        if self._mmap_region is not None:
            if handler_failed:
                # handler 未能干净 drain,可能仍有 in-flight DMA 引用 mmap 背书内存;
                # 强制同步设备,避免 unpin/unmap 与传输竞争。
                try:
                    torch.accelerator.synchronize()
                except Exception:
                    logger.warning(
                        "Device sync before mmap cleanup failed; "
                        "proceeding with cleanup anyway",
                        exc_info=True,
                    )
            self._mmap_region.cleanup()
            self._mmap_region = None

评论区精华

end_event.synchronize 失败后的处理方式 设计

orozery 给出 suggestion:同步失败时仅记录日志、清空队列并 break,不保留 sync_error 重抛,以简化控制流。

结论:未完全采纳:最终 head 版本仍保留 sync_error 并在 handler.shutdown() 末尾重新抛出,使其成为 worker 层感知失败、触发设备同步与统一清理的信号,属于后续 commit "harden CPU offload shutdown errors" 的加固内容。 · 已解决

test_transfer 中 del 引用块是否保留 正确性

orozery 转述 Claude 意见:不要删除 test_transfer 中的 del 块,因为测试局部引用在 worker.shutdown() 期间持有 buffer exports,删除会导致清理时序断言失真。

结论:采纳:最终 head 版本保留 del 块,仅移除手工 mmap_region.cleanup() 调用。 · 已解决

handler 失败后 cleanup 前的设备同步 正确性

orozery 建议在 handler 关停失败时先强制设备同步再执行 mmap_region.cleanup(),防止 mmap 卸载与 in-flight DMA 竞争。

结论:采纳:最终使用平台无关的 torch.accelerator.synchronize()(而非 torch.cuda.synchronize()),并带 try/except 警告降级路径。 · 已解决

风险与影响

  1. 构造失败路径未覆盖CPUOffloadingWorker.__init__self._mmap_region = mmap_region 之后若中途抛异常(如 pin_mmap_region、tensor 分配失败),region 不会进入 shutdown 流程,可能泄漏;该路径依赖 PR#51227 类的构造失败清理机制,需配套合入。
  2. 异常传播行为变化SingleDirectionOffloadingHandler.shutdown() 现在可能重新抛出 sync_error,若未来有其他调用方直接调用 handler 的 shutdown 且未捕获,行为与旧版(通常不抛)不同;当前仅 CPUOffloadingWorker.shutdown() 调用,已做捕获。
  3. torch.accelerator.synchronize() 依赖:该平台无关 API 在异常时走了 warning 降级路径,region 仍会清理,但极端情况下可能存在与 in-flight DMA 的竞争残留,属于可接受的防御降级。
  4. 测试真实性局限:新增测试大量依赖 MagicMock__new__ 注入,真实 CUDA event 的 device lost 等失败路径无法在 CI 复现,防护逻辑主要靠 code review 保障。

对用户:启用 CPU KV offloading 时,引擎 teardown 不再因事件同步失败而永久挂起,mmap 相关泄漏与崩溃概率降低(与 PR#51317 的硬退出治理形成互补)。对系统:关停顺序变得更可预测——store/load 两个方向均 drain 后统一释放 mmap region,错误通过异常日志可观测。对团队:确立了“单一所有权 + 分阶段清理 + 错误降级”的生命周期模式,测试中 __new__ + MagicMock 记录调用序列的写法可复用。影响范围仅限于 vllm/v1/kv_offload/cpu/ 的关停路径,正常推理与传输路径无行为变化。

关停路径重构 构造失败路径依赖配套 PR 错误路径测试依赖 Mock

关联 Issue

未识别关联 Issue

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

完整报告

参与讨论