执行摘要
- 一句话:异步 Mooncake 缓存查找,减少调度开销
- 推荐动作:值得精读,尤其关注
LookupKeyClient 如何利用 ThreadPoolExecutor 实现非阻塞 RPC,以及调度器通过返回 None 与重试机制配合的实现模式。对于需要将同步 I/O 异步化的类似场景,本 PR 提供了一个简洁的参考样板。
功能与动机
来自 PR body:
Looking up keys in Mooncake currently happens synchronously inside get_num_new_matched_tokens, on the scheduler's critical path. Each lookup costs ~1–2 ms per request on average, and that latency is paid serially during scheduling.
此变更为可选异步模式,默认保持同步行为,用户可按需启用以减少调度器阻塞。
实现拆解
-
LookupKeyClient 异步化 (worker.py):引入 ThreadPoolExecutor(max_workers=1) 和一个 futures 字典,将原来同步的 lookup 方法拆分为私有的 _lookup(实际 ZMQ 调用)和新的 lookup(提交并轮询)。新增 discard 方法取消并清理未完成的 Future。reset 也通过 executor 提交。所有 ZMQ 调用现在都发生在单一线程中,无需锁。
-
调度器集成 (scheduler.py):在 MooncakeStoreScheduler.__init__ 中解析 lookup_async 配置。get_num_new_matched_tokens 调用 client.lookup 时传入 non_block 标志。当异步查找未完成时返回 (None, False),V1 调度器会自动重试。在 build_connector_meta 中,对已完成的请求调用 client.discard() 以避免过时结果。
-
返回类型适配 (connector.py):接口返回类型从 tuple[int, bool] 拓宽为 tuple[int | None, bool],符合调度器重试契约。
-
测试覆盖:新增 3 个单元测试,使用 threading.Event 控制异步完成点,验证非阻塞 lookup、discard 行为以及调度器重试路径。更新调度器测试中的桩类以匹配新签名。
-
文档更新:在 mooncake_store_connector_usage.md 中添加 lookup_async 配置说明。
关键文件:
vllm/distributed/kv_transfer/kv_connector/v1/mooncake/store/worker.py(模块 存储 worker;类别 source;类型 core-logic;符号 lookup, _lookup, reset, discard): 核心变更文件,实现 LookupKeyClient 的异步化,包括 ThreadPoolExecutor、futures 字典、非阻塞 lookup 和 discard 逻辑。
vllm/distributed/kv_transfer/kv_connector/v1/mooncake/store/scheduler.py(模块 调度器;类别 source;类型 core-logic): 调度器集成异步查找逻辑,包括配置读取、lookup 调用修改和请求完成时 discard。
tests/v1/kv_connector/unit/test_mooncake_store_connector.py(模块 测试;类别 test;类型 test-coverage;符号 _poll_lookup, _gated_recv, recv, test_lookup_key_client_non_block_lookup_async): 新增 3 个测试覆盖异步查找、discard 和调度器重试场景。
vllm/distributed/kv_transfer/kv_connector/v1/mooncake/store/connector.py(模块 连接器;类别 source;类型 core-logic): 返回类型从 tuple[int, bool] 改为 tuple[int | None, bool],以传递推迟信号。
tests/v1/kv_connector/unit/test_mooncake_store_scheduler.py(模块 测试;类别 test;类型 test-coverage;符号 lookup): 更新 _StubLookupClient 签名以匹配新 lookup 接口,并设置 lookup_async 属性。
docs/features/mooncake_store_connector_usage.md(模块 文档;类别 docs;类型 documentation): 添加 lookup_async 配置项说明。
关键符号:LookupKeyClient.init, LookupKeyClient._lookup, LookupKeyClient.lookup, LookupKeyClient.discard, LookupKeyClient.reset, MooncakeStoreScheduler.init, MooncakeStoreScheduler.get_num_new_matched_tokens, MooncakeStoreScheduler.build_connector_meta
关键源码片段
vllm/distributed/kv_transfer/kv_connector/v1/mooncake/store/worker.py
核心变更文件,实现 LookupKeyClient 的异步化,包括 ThreadPoolExecutor、futures 字典、非阻塞 lookup 和 discard 逻辑。
# vllm/distributed/kv_transfer/kv_connector/v1/mooncake/store/worker.py
# LookupKeyClient 类,支持异步前缀缓存查找
class LookupKeyClient:
def __init__(self, vllm_config: VllmConfig):
self.encoder = MsgpackEncoder()
self.ctx = zmq.Context()
socket_path = get_zmq_rpc_path_lookup(vllm_config)
self.socket = make_zmq_socket(self.ctx, socket_path, zmq.REQ, bind=False)
# 使用单线程执行器处理所有 ZMQ 调用,避免锁
self.executor = ThreadPoolExecutor(
max_workers=1, thread_name_prefix='MooncakeLookupClient'
)
self.futures: dict[str, Future[int]] = {}
def _lookup(self, token_len: int, block_hashes: list[BlockHash]) -> int:
# 实际的 ZMQ 同步查找,在 executor 线程中运行
hash_strs = [h.hex() for h in block_hashes]
hash_frames = self.encoder.encode(hash_strs)
token_len_bytes = token_len.to_bytes(4, byteorder='big')
all_frames = [LOOKUP_MSG, token_len_bytes] + list(hash_frames)
self.socket.send_multipart(all_frames, copy=False)
resp = self.socket.recv()
return int.from_bytes(resp, 'big')
def lookup(
self,
req_id: str,
token_len: int,
block_hashes: list[BlockHash],
non_block: bool = False,
) -> int | None:
# 非阻塞查找:如果 non_block 为 True 且结果尚未就绪,返回 None
future = self.futures.get(req_id)
if future is None:
future = self.executor.submit(self._lookup, token_len, list(block_hashes))
self.futures[req_id] = future
if non_block and not future.done():
return None # 调度器重试
try:
return future.result()
except Exception as e:
logger.error('Async Mooncake lookup failed for %s: %s', req_id, e)
return 0
finally:
del self.futures[req_id]
def discard(self, req_id: str) -> None:
# 丢弃指定请求的查找任务,用于中止或已完成请求
future = self.futures.pop(req_id, None)
if future is not None:
future.cancel()
def _reset(self) -> bool:
# 实际的 reset 逻辑,在 executor 中执行
self.socket.send(RESET_MSG)
resp = self.socket.recv()
return bytes(resp) == RESP_OK
def reset(self) -> bool:
# 通过 executor 异步执行 reset
return self.executor.submit(self._reset).result()
def close(self):
self.executor.shutdown(wait=False, cancel_futures=True)
self.socket.close(linger=0)
vllm/distributed/kv_transfer/kv_connector/v1/mooncake/store/scheduler.py
调度器集成异步查找逻辑,包括配置读取、lookup 调用修改和请求完成时 discard。
# vllm/distributed/kv_transfer/kv_connector/v1/mooncake/store/scheduler.py
class MooncakeStoreScheduler:
def __init__(self, vllm_config: VllmConfig, kv_cache_config: KVCacheConfig):
kvc_extra_config = vllm_config.kv_transfer_config.kv_connector_extra_config
self.load_async = kvc_extra_config.get('load_async', True)
self.lookup_async = kvc_extra_config.get('lookup_async', False) # 新增:默认关闭
self.client = LookupKeyClient(vllm_config)
# ...
def get_num_new_matched_tokens(
self,
request: Request,
num_computed_tokens: int,
) -> tuple[int | None, bool]:
# 返回 (None, False) 表示异步查找未完成,调度器应重试
token_len = request.num_tokens // self._block_size * self._block_size
if token_len < self._block_size:
return 0, False
num_external_hit_tokens = self.client.lookup(
request.request_id,
token_len,
request.block_hashes,
non_block=self.lookup_async,
)
if num_external_hit_tokens is None:
return None, False # 查找尚未完成,让调度器后续重试
# 以下处理与之前相同(省略)
def build_connector_meta(self, scheduler_output: SchedulerOutput) -> KVConnectorMetadata:
force_skip_save = self.kv_role == 'kv_consumer'
for finished_req_id in scheduler_output.finished_req_ids:
self.client.discard(finished_req_id) # 丢弃已完成请求的查找 Future
# ... 其他清理 ...
评论区精华
- 函数命名讨论:reviewer zhewenl 建议将方法命名为
poll_lookup。作者 ivanium 认为 poll_lookup 隐含阻塞,更倾向于 try_lookup。最终实现中保留了 lookup 一词,通过 non_block 参数区分模式。
- 锁的必要性:Dao007forever 询问是否需要
socket_lock,ivanium 解释即使未来将 client 移到调度器侧,后台线程仍可能共享 socket,因此锁仍必要。但后续 njhill 提出了更简洁的 Futures 方案,完全避免了锁。
-
架构简化建议:njhill 在 review 中建议“使用单线程 ThreadPoolExecutor 和 Futures 来管理所有 socket 访问,避免使用锁”,并示范了提交。该方案被采纳,最终实现基于 njhill 的 commit,大幅精简了状态管理(去掉了自建的队列、线程和锁,改用 executor + futures 字典)。
-
方法命名:poll_lookup vs try_lookup (style): 最终实现保留了 lookup 名称,通过 non_block 参数区分阻塞与非阻塞模式。
- 锁的必要性讨论 (design): 最终采用 njhill 的 Futures 方案,完全避免了锁,使讨论 moot。
- 使用 Future 简化异步实现 (design): 该建议被采纳,最终实现基于 Futures 方案,代码简洁且线程安全。
风险与影响
- 风险:
- 默认关闭降低风险:
lookup_async 默认 False,现有用户行为不变,无回归风险。
- 异步一致性:若启用异步,调度器可能重试同一请求,需确保
discard() 在请求完成或中止时正确清理 Future,否则可能返回过时数据。当前实现通过 finally 块和 discard 中的 cancel() 保证。
- 线程模型:
ThreadPoolExecutor(max_workers=1) 确保所有 ZMQ 调用在单个线程中执行,避免锁竞争;但若查找请求大量积压,单线程吞吐可能成为瓶颈。由于每次查找通常 1-2ms,且与调度重叠,风险较低。
- 返回类型放宽:
tuple[int | None, bool] 的使用需要调用方(V1 调度器)已处理 None 分支,此 PR 依赖现成的 ext_tokens is None 重试机制,无需额外调度器改动。
- 影响:
- 用户影响:默认无感知;启用
lookup_async 后,调度器响应延迟降低(查找不再串行阻塞),但单请求完成时间可能略有增加(需多一次调度步骤重试)。适用于对延迟敏感的高并发前缀缓存场景。
- 系统影响:减少调度关键路径的不必要等待,提高整体吞吐。后台线程执行 ZMQ 网络 I/O,对 GPU 计算无干扰。
- 团队影响:新的配置项需要文档和版本兼容说明;后续若将
LookupKeyClient 移入调度器进程,此异步机制仍然是安全的。
- 风险标记:默认关闭降低风险, 异步清理保证, 单线程 executor 可能瓶颈
关联脉络
- PR #45371 [Bugfix][KV Connector] Disable Mooncake TP put-striding when DCP > 1: 修改了相同的 Mooncake store worker 文件(worker.py),属于同一套件的基础设施改进。
参与讨论