Prhub

#45971 [Perf][KVConnector][Mooncake] Parallelize KV load with a receive-thread pool

原始 PR 作者 ivanium 合并时间 2026-06-25 09:25 文件变更 3 提交数 3 评论 2 代码增减 +54 / -27

执行摘要

并行化 Mooncake KV 加载接收线程

MooncakeStoreConnector 中每个 KV 加载请求在 RDMA 传输开始前需要支付 Python 端请求准备和 master key 查找的控制开销。单接收线程下这些开销是串行的,导致 NIC 在连续传输之间空闲。PR body 指出:“the RDMA bw usage is bottlenecked for long sequences”。

该 PR 设计清晰,变更范围小,风险低,值得合并。建议读者关注:共享队列的设计如何实现负载均衡,以及完成结果汇聚的线程安全处理。对于部署 Mooncake 的用户,建议实测不同线程数下的性能差异。

讨论亮点

Review 中仅有一次交互:zhewenl 询问为何移除 assert self.kv_recv_thread is not None,认为可以保留以防止线程未启动。ivanium 回应原断言主要用于通过 ruff 类型检查,并非功能必要。最终该断言被移除,因为线程池的启动已在初始化时确保。

实现拆解

  1. 环境变量注册:在 vllm/envs.py 中新增 VLLM_MOONCAKE_LOAD_RECV_THREADS(默认 1),控制接收线程数。
  2. 线程池创建MooncakeStoreWorker.__init__self.kv_recv_thread 替换为 self.kv_recv_threads: list[KVCacheStoreRecvingThread],并根据环境变量创建指定数量的线程,所有线程共享同一个请求队列 self.recv_request_queue
  3. 共享队列分发KVTransferThread.__init__ 增加可选的 request_queue 参数,允许外部传入共享队列;MooncakeStoreWorker.get_finished() 将请求直接放入共享队列,不再调用单个线程的 add_request
  4. 完成结果汇聚get_finished()get_block_ids_with_load_errors() 遍历所有接收线程,汇聚各线程的完成请求 ID 和加载错误块 ID。
  5. 测试适配:在 tests/v1/kv_connector/unit/test_mooncake_store_worker.py 中,将 w.kv_recv_thread 改为 w.kv_recv_threads 列表,并初始化共享队列。
文件 模块 状态 重要度
vllm/distributed/kv_transfer/kv_connector/v1/mooncake/store/worker.py 分布式 modified 6.82
vllm/envs.py 配置 modified 4.91
tests/v1/kv_connector/unit/test_mooncake_store_worker.py 测试 modified 3.93

关键符号

MooncakeStoreWorker.register_kv_caches MooncakeStoreWorker.get_finished MooncakeStoreWorker.get_block_ids_with_load_errors KVCacheStoreRecvingThread.__init__

关键源码片段

vllm/distributed/kv_transfer/kv_connector/v1/mooncake/store/worker.py core-logic

核心实现文件,实现了接收线程池的创建、请求分发、结果汇聚。

# vllm/distributed/kv_transfer/kv_connector/v1/mooncake/store/worker.py
# 关键变更:线程池创建与共享队列class MooncakeStoreWorker:
    def __init__(self, ...):
        # ...
        # 以前 : self.kv_recv_thread = None
        # 现在 : 使用线程池
        self.kv_recv_threads: list[KVCacheStoreRecvingThread] = []
        self.num_recv_threads = max(1, envs.VLLM_MOONCAKE_LOAD_RECV_THREADS)
        # 所有接收线程共享一个请求队列
        self.recv_request_queue: queue.Queue[ReqMeta] = queue.Queue()
​
    def register_kv_caches(self, ...):
        # ...
        # 创建线程池,所有线程从同一个队列取任务
        self.kv_recv_threads = []
        for i in range(self.num_recv_threads):
            ready_event = threading.Event()
            recv_thread = KVCacheStoreRecvingThread(
                self.store,
                self.coord,
                self.token_dbs,
                self.block_size,
                self.tp_rank,
                ready_event,
                disk_offload_buffer_budget_bytes=...,
                record_operation=...,
                request_queue=self.recv_request_queue, # 传递共享队列
            )
            recv_thread.name = f"KVCacheStoreRecvingThread-{i}"
            recv_thread.start()
            self.kv_recv_threads.append(recv_thread)
            ready_events_recving.append(ready_event)
        # 等待所有线程就绪
        for e in ready_events_recving:
            e.wait()
​
    def get_finished(self):
        # 汇聚所有线程的完成请求 ID
        finished: set[str] = set()
        for thread in self.kv_recv_threads:
            finished |= thread.get_and_clear_finished_requests()
        # ... 处理 kv events ...
        return finished
​
    def get_block_ids_with_load_errors(self):
        # 汇聚所有线程的加载错误块 ID
        result: set[int] = set()
        for thread in self.kv_recv_threads:
            result |= thread.get_and_clear_block_ids_with_load_errors()
        return result
vllm/envs.py core-logic

注册新的环境变量以控制接收线程数。

# vllm/envs.py
# 新增环境变量定义(类型注解)
VLLM_MOONCAKE_LOAD_RECV_THREADS: int = 1# 新增环境变量读取逻辑
"VLLM_MOONCAKE_LOAD_RECV_THREADS": lambda: int(
    os.getenv("VLLM_MOONCAKE_LOAD_RECV_THREADS", "1")
),

评论区精华

移除 assert self.kv_recv_thread is not None question

zhewenl 询问为何移除断言,认为可保留以防止线程未启动。ivanium 回应断言主要用于通过 ruff 类型检查,非功能必要。

结论:移除断言,因为线程池的创建在初始化时已确保。 · 已解决

风险与影响

主要风险在于线程安全:多个线程共享同一个 recv_request_queue,但 Python 的 queue.Queue 是线程安全的,因此风险较低。get_finished()get_block_ids_with_load_errors() 遍历线程列表时,没有对列表进行加锁保护,但线程创建后列表不再修改,因此问题不大。默认线程数为 1,不改变现有行为,回归风险低。

对用户:默认无行为变化,需要显式设置环境变量 VLLM_MOONCAKE_LOAD_RECV_THREADS 以启用多线程加速。对系统:在长序列或小批量加载场景下可显著提升 RDMA 带宽利用率。对团队:引入的新配置项需要文档说明和性能调优指导。

多线程共享队列 默认行为无变化

关联 Issue

未识别关联 Issue

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

完整报告

参与讨论