Prhub

#46278 [Bugfix][KVConnector] Fix SimpleCPUOffloadConnector GPU->CPU store race

原始 PR 作者 Saddss 合并时间 2026-06-23 04:08 文件变更 4 提交数 4 评论 16 代码增减 +244 / -16

执行摘要

修复 CPU 卸载存储与计算流的竞争条件

源自 Issue #45704:在生产负载下,SimpleCPUOffloadConnector 的 GPU->CPU 存储可能读取部分写入或过时的 KV 块,导致 CPU 缓存静默损坏,后续请求输出乱码。根本原因是 v1 重叠执行下主机先于 GPU 推进,而存储在独立的复制流上执行且缺少跨流同步;加上 srcAccessOrder=ANY 允许 DMA 在源数据未就绪时开始读取,两者结合造成了竞争。

此 PR 值得精读,尤其对于涉及 GPU 异步传输、跨流同步的开发者。它展示了一个典型的 CUDA 流间竞争 bug 的分析与修复过程,以及通过自验证测试确保修复的可靠性。对 srcAccessOrder 的讨论提供了深层理解 CUDA DMA 语义的机会。

讨论亮点
  • srcAccessOrder 必要性(ivanium vs Saddss):ivanium 最初认为仅事件同步即可修复,Saddss 引用 CUDA 驱动文档并解释:ANY 允许 DMA 在流屏障事件之前预读源地址,因此即使计算流完成,DMA 也可能读取未写入的内存;STREAM 强制 DMA 遵循流顺序,是安全保证所必需。ivanium 最终接受。
  • 注释精简(ivanium):ivanium 指出多处注释过于冗长,带有 agent 风格,建议简化。Saddss 逐步将核心 rationale 集中在文档字符串和枚举定义处,移除了内联重复说明。
  • 事件复用(ivanium):ivanium 建议将每次存储新分配的事件改为复用单个事件,只重新记录以减少开销。Saddss 采纳,引入 self._store_compute_done 成员变量。

实现拆解

  1. 定义常量与扩展签名:在 cuda_mem_ops.py 中添加 CU_MEMCPY_SRC_ACCESS_ORDER_STREAM(1)和 ANY(3)常量,替代硬编码魔数;build_params 新增 src_access_order 形参,默认 ANY
  2. 修改复制后端:在 copy_backend.pyDmaCopyBackend.init 中,存储参数使用 STREAM、加载参数使用 ANYlaunch_copy 增加 wait_event 参数并传递到队列;_copy_loop 在调用 copy_blocks 前根据 wait_event 让目标流等待事件完成。
  3. 修改工作器:在 worker.pySimpleCPUOffloadWorker 中增加 _store_compute_done 事件成员,get_finished 中每次存储前记录当前计算流的事件,作为 wait_event 传入 launch_copy;加载路径不变。
  4. 新增测试:新建 tests/v1/simple_kv_offload/test_worker.py,包含三个测试用例:test_store_orders_after_compute_write 使用 torch.cuda._sleep 制造确定性竞争,断言无屏障时损坏 >0、有屏障时损坏=0;test_get_finished_passes_wait_event_for_store_only 通过桩后端验证事件传递;test_build_params_src_access_order 验证参数覆盖。
文件 模块 状态 重要度
vllm/v1/simple_kv_offload/copy_backend.py 复制后端 modified 6.97
vllm/v1/simple_kv_offload/worker.py 工作器 modified 6.39
vllm/v1/simple_kv_offload/cuda_mem_ops.py CUDA 内存 modified 5.89
tests/v1/simple_kv_offload/test_worker.py 测试 added 7.66

关键符号

DmaCopyBackend.init DmaCopyBackend.launch_copy DmaCopyBackend._copy_loop SimpleCPUOffloadWorker.__init__ SimpleCPUOffloadWorker.get_finished build_params test_store_orders_after_compute_write test_get_finished_passes_wait_event_for_store_only test_build_params_src_access_order

关键源码片段

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

核心逻辑变更:在 `init` 中为存储参数指定 `srcAccessOrder=STREAM`,`launch_copy` 新增 `wait_event` 参数,`_copy_loop` 在复制前等待事件。

# vllm/v1/simple_kv_offload/copy_backend.py
from __future__ import annotations
import queue, threading
import torch
from vllm.v1.simple_kv_offload.cuda_mem_ops import (
    CU_MEMCPY_SRC_ACCESS_ORDER_ANY,
    CU_MEMCPY_SRC_ACCESS_ORDER_STREAM,
    BatchMemcpyParams,
    build_params,
    copy_blocks,
)class DmaCopyBackend:
    def init(
        self,
        gpu_caches: dict[str, torch.Tensor],
        cpu_caches: dict[str, torch.Tensor],
        device: torch.device,
        load_stream: torch.cuda.Stream,
        store_stream: torch.cuda.Stream,
    ) -> None:
        self._load_stream = load_stream
        self._store_stream = store_stream
        # 存储源是活跃的 KV 缓存 -> STREAM(与 get_finished 中的计算完成事件配合)
        self._store_params = build_params(
            gpu_caches, cpu_caches, store_stream,
            src_access_order=CU_MEMCPY_SRC_ACCESS_ORDER_STREAM,
        )
        # 加载源是稳定的固定主机内存 -> ANY
        self._load_params = build_params(
            cpu_caches, gpu_caches, load_stream,
            src_access_order=CU_MEMCPY_SRC_ACCESS_ORDER_ANY,
        )
        # ...(队列和线程设置省略)
​
    def launch_copy(
        self,
        src_blocks, dst_blocks, is_store, event_idx, events_list,
        wait_event=None, # 新增参数
    ) -> None:
        params = self._store_params if is_store else self._load_params
        self._queue.put(
            (src_blocks, dst_blocks, params, is_store,
             event_idx, events_list, wait_event)
        )
​
    @staticmethod
    def _copy_loop(q, device, load_stream, store_stream):
        current_platform.set_device(device)
        while True:
            item = q.get()
            if item is None:
                return
            (src_blocks, dst_blocks, params, is_store,
             event_idx, events_list, wait_event) = item
            stream = store_stream if is_store else load_stream
            if wait_event is not None:
                stream.wait_event(wait_event) # 确保计算完成后才复制
            copy_blocks(src_blocks, dst_blocks, params)
            event = torch.Event()
            event.record(stream)
            events_list.append((event_idx, event))
vllm/v1/simple_kv_offload/worker.py core-logic

调用端:`get_finished` 中记录计算完成事件并传递给 `launch_copy`,同时为存储事件提供生命周期管理。

# vllm/v1/simple_kv_offload/worker.py
class SimpleCPUOffloadWorker:
    def __init__(self, ...):
        # ...
        # 复用的事件对象,每次步骤重新记录
        self._store_compute_done: torch.Event | None = None
​
    def get_finished(self, finished_req_ids):
        """提交传输并报告已完成的事件。        存储(GPU->CPU)读取活跃的 KV 缓存,该缓存可能仍在被计算流写入
        (在 v1 重叠执行下),因此需要在当前流上记录计算完成事件并等待。
        加载(CPU->GPU)读取稳定的固定主机内存,立即启动。
        """
        metadata = self._connector_metadata
        if metadata is not None:
            # 加载:无需事件屏障
            if metadata.load_cpu_blocks:
                self._backend.launch_copy(
                    metadata.load_cpu_blocks, metadata.load_gpu_blocks,
                    is_store=False, event_idx=metadata.load_event,
                    events_list=self._load_events)
            # 存储:记录计算完成事件并传递
            if metadata.store_gpu_blocks:
                if self._store_compute_done is None:
                    self._store_compute_done = torch.Event()
                self._store_compute_done.record(torch.cuda.current_stream())
                self._backend.launch_copy(
                    metadata.store_gpu_blocks, metadata.store_cpu_blocks,
                    is_store=True, event_idx=metadata.store_event,
                    events_list=self._store_events,
                    wait_event=self._store_compute_done)
        # ...(其余逻辑)
vllm/v1/simple_kv_offload/cuda_mem_ops.py core-logic

定义了 `srcAccessOrder` 常量并扩展了 `build_params` 签名,为差异化访问顺序提供基础设施。

# vllm/v1/simple_kv_offload/cuda_mem_ops.py
# CUmemcpySrcAccessOrder 值(CUDA 驱动 API)
CU_MEMCPY_SRC_ACCESS_ORDER_STREAM = 1 # 安全用于仍在写入的源
CU_MEMCPY_SRC_ACCESS_ORDER_ANY = 3 # 仅安全用于稳定的源def build_params(
    src_caches: dict[str, torch.Tensor],
    dst_caches: dict[str, torch.Tensor],
    stream: torch.cuda.Stream,
    src_access_order: int = CU_MEMCPY_SRC_ACCESS_ORDER_ANY, # 新增参数
) -> BatchMemcpyParams:
    # ...(检查和构建数组)
    attrs = _CUmemcpyAttributes(srcAccessOrder=src_access_order)
    return BatchMemcpyParams(
        # ...
        attrs=attrs,
        # ...
    )

评论区精华

srcAccessOrder 必要性 正确性

ivanium 询问仅事件同步是否足以修复,srcAccessOrder=STREAM 是否必要。Saddss 引用 CUDA 驱动文档,解释 ANY 允许 DMA 在流屏障事件之前预读源地址,即使计算流完成,DMA 仍可能读取未写入的内存;STREAM 强制遵循流顺序,是安全保证所必需。

结论:ivanium 接受解释,认为设置 STREAM 是必要的。 · 已解决

精简注释 style

ivanium 指出多处注释过于冗长,带有 agent 风格,建议简化。Saddss 逐步将核心 rationale 集中在文档字符串和枚举定义处,移除了内联重复说明。

结论:注释已按建议精简。 · 已解决

复用计算完成事件 性能

ivanium 建议将每次存储时新分配的事件改为复用单个事件,只重新记录以减少开销。

结论:Saddss 采纳,引入 self._store_compute_done 成员变量并重用。 · 已解决

风险与影响

  • 事件同步逻辑如果错误(如事件记录在错误流上)可能导致存储仍与计算竞争或无限阻塞,但测试覆盖了正向(有屏障)和反向(无屏障)场景,且通过 150 次确定性运行验证。
  • build_params 增加了可选参数,但该函数目前只在 DmaCopyBackend 内部调用,无外部消费者,接口兼容。
  • 新测试使用 torch.cuda._sleep 这一非公共 API,不同 CUDA 版本可能行为不同,但测试通过即证明逻辑有效;临时替代方案(如 spin kernel)会增加复杂度。
  • 用户:修复了启用 CPU 卸载时间歇性输出乱码的问题,提升推理可靠性;性能 A/B 测试显示吞吐变化 -0.47%、TPOT 变化 -0.1%,在噪声范围内。
  • 系统:仅影响 SimpleCPUOffloadConnector 路径,对默认无卸载或其他连接器无影响;存储事件和 srcAccessOrder 变更仅增加一次 event record 和 wait,开销极低。
  • 团队:为异步卸载的正确性提供了明确模式(计算完成事件 + 流顺序),可供后续其他连接器参考。
GPU- 流同步 CUDA 特有 非公共 API 测试

关联 Issue

#45704 [Bug]: SimpleCPUOffloadConnector corrupts restored KV under load → intermittent garbled output (missing GPU→CPU store cross-stream sync)

完整报告

参与讨论