Prhub

#41383 [Nixl][PD] Lease renewal TTL KV blocks on P

原始 PR 作者 NickLucche 合并时间 2026-05-11 17:27 文件变更 13 提交数 17 评论 51 代码增减 +462 / -48

执行摘要

新增 NIXL 心跳租约续期,优化 KV 块保留时间

原有的 VLLM_NIXL_ABORT_REQUEST_TIMEOUT 单一超时过于简单,当 D 崩溃时会导致 P 端长时间持有远程请求的 KV 块,造成资源浪费。需要一套动态 TTL "租约续期" 系统,最小化块在 P 端的滞留时间,同时让 D 能够在等待队列拥堵时延长 P 端块的租约。设计文档详见 Google Docs 链接。

值得所有关注分布式 KV 传输的开发者精读。设计上展示了典型的租约续期模式,接口抽象和配置简化决策值得借鉴。建议重点阅读 scheduler.pyon_new_request/_stop_heartbeatworker.py_ensure_handshake/_send_heartbeats,体会如何在分布式组件间实现轻量级心跳协议。

讨论亮点

主要讨论集中在接口设计和配置简化。

  • @markmc 建议将三个配置项合并为单一 kv_lease_duration,避免用户暴露过多细节。作者采纳,最终接口仅暴露 kv_lease_duration
  • @orozery 提议调度器只需在 add_request 时通知连接器,避免复杂的 ObservableQueue 抽象。经过讨论,作者放弃 ObservableRequestQueue,改用直接的 connector.on_new_request 回调。
  • @ivanium 质疑 ObservableRequestQueue 的放置位置,建议移到连接器内部,最终决定保持简单。
  • @markmc 指出 _send_heartbeats 访问 self._remote_agents 时未持有 _handshake_lock,存在数据竞争。作者修复。
  • 关于双向传输场景,@markmc 认为 D 端块的 TTL 应独立配置且与心跳无关,最终引入 decoder_kv_blocks_ttl 配置,默认 480 秒。

实现拆解

  1. 元数据扩展:新增 HeartbeatInfo dataclass,在 NixlConnectorMetadata 中添加 heartbeat_by_engine 字典,用于调度器传递心跳信息给 Worker。
  2. 调度器心跳追踪:在 NixlConnectorScheduler 中增加 _heartbeat_by_engine_heartbeat_req_engine_last_heartbeat_time 数据结构,实现 on_new_request 方法在请求入队时注册远程引擎关联的心跳信息,_stop_heartbeat 方法在请求发送完成或终止时移除对应条目。心跳发送节流通过比较 _last_heartbeat_time 与当前时间控制。
  3. Worker 心跳收发:新增 _send_heartbeats 方法,在 start_load_kv 时从元数据中提取心跳信息并向每个远程引擎发送心跳消息(ZMQ);_handle_heartbeat 在 P 端接收心跳后延长 _reqs_to_send 中对应请求的过期时间(time.perf_counter() + self._lease_extension)。
  4. 接口整合KVConnector 基类增加 on_new_requestupdate_connector_output 方法,分别用于通知连接器新请求入队和处理 Worker 输出(停止心跳)。NixlConnectorMultiConnector 转发调用。
  5. 配置简化:将原先的三个配置项(heartbeat_intervalheartbeat_lease_extensioninitial_kv_lease)合并为单一 kv_lease_duration,内部派生 _heartbeat_interval = kv_lease_duration // 6_lease_extension = kv_lease_duration * 2 // 3。同时新增 decoder_kv_blocks_ttl 配置用于双向传输中 D 端块的固定 TTL。
  6. 测试覆盖:新增 tests/v1/kv_connector/unit/test_nixl_heartbeat.py,包含 5 个单元测试验证心跳跟踪、分组、节流、停止路径。同时修改测试工具函数 make_nixl_scheduler 支持心跳参数。
文件 模块 状态 重要度
vllm/distributed/kv_transfer/kv_connector/v1/nixl/scheduler.py NIXL 调度器 modified 8.5
vllm/distributed/kv_transfer/kv_connector/v1/nixl/worker.py NIXL Worker modified 8.44
tests/v1/kv_connector/unit/test_nixl_heartbeat.py 心跳测试 added 7.62
vllm/distributed/kv_transfer/kv_connector/v1/nixl/metadata.py 元数据 modified 6.48
vllm/distributed/kv_transfer/kv_connector/v1/nixl/connector.py 连接器桥接 modified 6.34
vllm/distributed/kv_transfer/kv_connector/v1/base.py 基类接口 modified 5.5

关键符号

on_new_request _stop_heartbeat update_connector_output _handle_heartbeat _ensure_handshake _send_heartbeats

关键源码片段

vllm/distributed/kv_transfer/kv_connector/v1/nixl/worker.py dependency-wiring

Worker 端实现心跳收发和租约扩展,重构 handshake 逻辑

    def _ensure_handshake(
        self,
        engine_id: EngineId,
        host: str,
        port: int,
        tp_size: int,
    ) -> Future[dict[int, str]] | None:
        """
        确保与远程引擎的握手已完成或正在执行。        如果握手已成功完成则直接返回 ``None``,如果握手正在进行则返回 ``Future``,
        如果尚未开始则启动并返回 ``Future``。
        调用方可以在返回的 ``Future`` 上附加 per-request 回调。
        """
        with self._handshake_lock:
            # 检查是否已有成功的远程代理连接
            if engine_id in self._remote_agents:
                return None # 握手已完成
            # 检查是否有正在进行的握手
            fut = self._handshake_futures.get(engine_id)
            if fut is not None:
                return fut # 等待中的握手
            # 提交新的握手任务
            fut = self._handshake_initiation_executor.submit(
                self._nixl_handshake,
                host,
                port,
                tp_size,
                engine_id,
            )
            self._handshake_futures[engine_id] = fut
​
            def done_callback(f: Future[dict[int, str]], eid=engine_id):
                """握手完成后的清理与错误日志回调"""
                with self._handshake_lock:
                    del self._handshake_futures[eid]
                try:
                    f.result() # 触发异常如果失败
                except Exception as e:
                    # 记录握手失败日志
                    logger.error('Handshake failed for engine %s: %s', eid, e)
​
            fut.add_done_callback(done_callback)
            return fut
tests/v1/kv_connector/unit/test_nixl_heartbeat.py test-coverage

新增单元测试覆盖心跳功能,验证追踪、节流、停止路径

def test_on_new_request_tracks_and_groups():
    """添加两个请求到同一引擎,一个到另一引擎;验证分组正确"""
    s = _sched() # 创建带心跳支持的调度器
    s.on_new_request(_req(1)) # 添加请求 1(默认引擎 _ENGINE_A)
    s.on_new_request(_req(2)) # 添加请求 2(同一引擎)
​
    # 验证同一引擎下两个请求的心跳 ID 集合
    assert s._heartbeat_by_engine[_ENGINE_A].req_ids == {'prefill-1', 'prefill-2'}
    info = s._heartbeat_by_engine[_ENGINE_A]
    # 验证远程主机、端口、TP 大小与模拟数据一致
    assert (info.host, info.port, info.tp_size) == ('my-host', 1234, 1)
    # 验证反向映射表
    assert s._heartbeat_req_engine['id-1'] == (_ENGINE_A, 'prefill-1')
​
    # 添加不同引擎的请求
    r3 = _req(3)
    r3.kv_transfer_params['remote_engine_id'] = 'engine-b'
    s.on_new_request(r3)
    # 此时应有两个引擎条目
    assert len(s._heartbeat_by_engine) == 2

评论区精华

简化配置项数量 设计

@markmc 认为三个配置项(heartbeat_interval、heartbeat_lease_extension、initial_kv_lease)过于繁琐,建议合并为一个 kv_lease_duration。

结论:采纳,最终只暴露 kv_lease_duration 给用户,内部按固定比例派生其他值。 · 已解决

远程引擎追踪优化 性能

@markmc 建议在 on_new_request 中避免为每个新请求创建 HeartbeatInfo,因为远程引擎出现频率低,应先检查再创建。

结论:采纳,修改为仅在 remote_engine_id 不在字典中时创建新 HeartbeatInfo。 · 已解决

心跳竞态条件 正确性

@gemini-code-assist[bot] 指出 _handle_heartbeat 中更新 _reqs_to_send 的过期时间可能存在竞态,建议使用线程安全更新。

结论:作者认为问题不大,未做额外修改。markmc 也认为不大可能覆盖更晚的过期时间。状态 unresolved。 · unresolved

ObservableRequestQueue 设计争议 设计

@ivanium 质疑 ObservableRequestQueue 放在调度器中的必要性,提议移到连接器内部。最终同意 @orozery 建议,保持简单,移除 ObservableQueue。

结论:移除 ObservableRequestQueue,改用直接的 connector.on_new_request 回调。 · 已解决

双向传输模式 TTL 独立配置 设计

@markmc 认为双向传输下 D 端块的 TTL 应独立配置,且与心跳无关,建议使用 decoder_kv_blocks_ttl。

结论:采纳,新增 decoder_kv_blocks_ttl 配置,默认为 480 秒。 · 已解决

锁保护 remote_agents 访问 安全

@markmc 指出 _send_heartbeats 访问 self._remote_agents 时未持有 _handshake_lock,可能导致数据竞争。

结论:作者修复,添加 with self._handshake_lock 保护。 · 已解决

风险与影响

  1. 线程安全_send_heartbeats_ensure_handshake_remote_agents 的访问虽已加锁,但其他路径(如 _handle_heartbeat 更新 _reqs_to_send)未显式加锁,可能在高并发下出现竞态。
  2. 配置兼容性:废弃 VLLM_NIXL_ABORT_REQUEST_TIMEOUT 环境变量,用户需迁移到新的 kv_lease_duration 配置,否则将使用默认值 30 秒。
  3. 双向传输风险:P 端不发送心跳,D 端依赖固定 TTL,若 decoder_kv_blocks_ttl 设置过短,在 D 端请求处理延迟时可能导致块提前释放。
  4. 性能风险:心跳消息增加网络开销,但节流机制(每 kv_lease_duration/6 秒一次)将其控制在可接受范围。

影响范围限于使用 NIXL KV 连接器的分布式推理场景。主要好处是降低了 P 端块的无谓滞留,提高显存利用率,并允许 D 端通过心跳主动延长块租约,减少因拥塞导致的块提前释放。用户需要为 kv_lease_duration 设置合适的值。测试发现对正常吞吐无显著影响。团队需维护新增的 4 个配置项和对应的测试。

核心路径变更 线程安全 双向传输兼容性 配置废弃

关联 Issue

未识别关联 Issue

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

完整报告

参与讨论