Files
tech/知识/微信公众平台/marcus/2026-08-05_vLLM_PD分离_源码级机制_marcus.md
T
arno c0ba3fb853 文档(金鹏): 2026-08-06 章节 48 篇文章摘要归档
- 46 篇原文+摘要双文件归档(按 来源/作者 分层,复用本地归档 20 篇+新抓取 26 篇)
- 即梦生成 9 组主题配图(大图+列表缩略图)存入 知识/金鹏/20260806/
- 章节重组为 9 个主题分组并挂接摘要引用
2026-08-06 18:00:50 +08:00

13 KiB
Raw Blame History

vLLM PD 分离

来源:微信公众平台
作者marcus
发布日期:2026-08-05(原文未标注日期,按下载日占位)
原文链接https://mp.weixin.qq.com/s/p6gh4hK25PujQ1UJzFTbLg

本文是 vLLM 基于 Mooncake Connector 实现 PD 分离的源码级机制剖析,覆盖 Proxy 层、请求传递、Scheduler 调度(P producer / D consumer)、Worker 层 RDMA KV Cache 传输、传输后恢复执行、完整时序与关键设计要点。


一、Proxy 层:同一份 Prompt 同时发给 P 和 D

是的,Proxy 会把原始 prompt 同时发给 Prefill 和 Decode 节点,但携带不同的 kv_transfer_params 标志:

Client Prompt "Hello, how are you?"  
         │  
         ▼  
    mooncake_connector_proxy.py  
         │  
    ┌────┴────┐  
    ▼         ▼  
 Prefill    Decode  
 节点        节点

发给 Prefill 节点send_request_to_service, 第 250280 行):

req_data["kv_transfer_params"] = {  
    "do_remote_decode": True,      # ← "我是 P,算完把 KV 传出去"  
    "do_remote_prefill": False,  
    "transfer_id": f"xfer-{request_id}",  
}  
req_data["max_tokens"] = 1         # ← 只算 1 个 token 就停  
req_data["stream"] = False         # ← 不等流式返回  

asyncio.create_task fire-and-forget 发送,不等结果。

发给 Decode 节点stream_service_response, 第 283312 行):

req_data["kv_transfer_params"] = {  
    "do_remote_decode": False,  
    "do_remote_prefill": True,     # ← "我是 D,从远端的 P 拉 KV Cache"  
    "remote_bootstrap_addr": ...,  # ← P 的 bootstrap 地址  
    "remote_engine_id": ...,       # ← P 的 Mooncake engine ID(含 DP rank  
    "transfer_id": f"xfer-{request_id}",  
}  

以流式方式等待返回,逐 chunk 回传客户端。

每个请求都有唯一的 transfer_idUUID——这是 P 和 D 两端关联同一次 KV Cache 传输的关键标识。


二、请求进入引擎:kv_transfer_params 的传递路径

Proxy → API Server (SamplingParams.extra_args)  
      → EngineCoreRequest  
      → Request.__init__() 提取 kv_transfer_params  
      → Scheduler 调度  
      → MooncakeConnector 处理  

vllm/v1/request.py 第 116119 行:

if sampling_params.extra_args is not None:  
    self.kv_transfer_params = sampling_params.extra_args.get(  
        "kv_transfer_params"  
    )  

kv_transfer_params 作为 Request 对象的属性一直保留,供 Scheduler 和 KVConnector 使用。


三、Scheduler 调度层:区分本地计算和远程 KV 加载

这是 PD 分离最核心的逻辑,在 vllm/v1/core/sched/scheduler.pyschedule() 方法中。

3.1 Prefill 节点 (is_kv_producer) 的调度

do_remote_decode=True 时:

  1. get_num_new_matched_tokensMooncakeConnectorScheduler 第 712755 行):返回 (0, False)——P 节点不需要加载任何外部 KV,正常计算 prompt。
  2. update_state_after_alloc(第 757–805 行):记录该请求需要将 KV Cache 发送出去:
    elif params.get("do_remote_decode"):
        self._reqs_need_send[request.request_id] = (request, [])
    
  3. request_finished(第 838–892 行):当请求完成时(只生成了 1 个 token),不立即释放 block,而是标记为需要异步发送:
    if delay_free_blocks:
        self._reqs_need_send[request.request_id] = (
            request, self.get_sw_clipped_blocks(block_ids),
        )
        return delay_free_blocks, None  # ← True = 延迟释放
    

3.2 Decode 节点 (is_kv_consumer) 的调度

do_remote_prefill=True 时:

  1. get_num_new_matched_tokens:返回 (num_prompt_tokens, True)——告诉 Scheduler"这 N 个 token 的 KV 已经在远端算好了,异步去拉"。
    if params.get("do_remote_prefill"):
        token_ids = request.prompt_token_ids or []
        count = self._get_remote_prefill_token_count(len(token_ids)) - (
            num_computed_tokens
        )
        if count > 0:
            return count, True  # ← count 是 prompt 长度,True = 异步加载
    
  2. Scheduler 分配 KV block(第 953971 行):allocate_slots 中传入 num_external_computed_tokens,并为远程 token 预留 block。关键参数:
    • delay_cache_blocks=True — block 分配了但不立即写入 prefix cache(因为数据还在远端)
    • reserved_blocks — 确保有足够空间完成整个异步传输,防止死锁
  3. 请求进入 WAITING_FOR_REMOTE_KVS 状态(第 10101040 行):
    if load_kv_async:
        request.status = RequestStatus.WAITING_FOR_REMOTE_KVS
        step_skipped_waiting.prepend_request(request)
        request.num_computed_tokens = num_computed_tokens
        self._inflight_prefills.add(request)
        # 跳过 zeroing:异步写入可能与 zeroing 竞态
    
    此时 request 被放入 skipped_waiting 队列,不会进入 running 列表,本轮不参与 GPU 计算。
  4. build_connector_meta:将需要 recv 的请求打包成 MooncakeConnectorMetadata(含 remote_engine_idremote_bootstrap_addrtransfer_idlocal_block_ids),放入 SchedulerOutput.kv_connector_metadata 传递给 Worker。

四、Worker 层:Mooncake 的 KV Cache RDMA 传输

4.1 D 节点 Worker:拉取 KV Cache

start_load_kvMooncakeConnectorWorker, 第 20282039 行):

def start_load_kv(self, metadata: MooncakeConnectorMetadata):  
    if not self.is_kv_producer and metadata.reqs_to_recv:  
        asyncio.run_coroutine_threadsafe(  
            self._start_load_kv(metadata.reqs_to_recv), self.receiver_loop  
        )  

_start_load_kvhandle_new_engine_idreceive_kvreceive_kv_from_single_worker

  1. 首先通过 _connect_to_prefiller_bootstrap(remote_bootstrap_addr) 查询 P 节点的 bootstrap server,获取 P 侧所有 TP/PP Worker 的 ZMQ 地址
  2. 构建 MooncakeXferMetadata(第 1825–1840 行),包含:
    • remote_hostname/remote_port:D 侧的地址,P 侧用来做 RDMA 写入
    • remote_tp_size/remote_tp_rankTP 拓扑信息
    • req_blocks{req_id: (transfer_id, local_block_ids)}——D 侧已分配的 block ID
    • kv_caches_base_addr/block_lens/kv_block_lens:D 侧 KV Cache 的物理内存地址和布局
  3. 通过 ZMQ DEALER socket 向 P 侧对应 TP rank 的 Worker 发送此元数据
  4. P 侧收到后,通过 RDMA (batch_transfer_sync_write) 将 KV Cache 数据直接写入 D 侧的 GPU 内存地址
  5. 传输完成后,D 侧标记 finished_recving_reqs

4.2 P 节点 Worker:推送 KV Cache

record_send_reqs(第 20002026 行):P 侧的 Worker 将 Scheduler 传递来的 reqs_to_send(含 transfer_idlocal_block_ids)记录到 self.reqs_need_send 字典。

send_kv_to_decode(第 11931382 行):当 D 侧的 ZMQ 请求到达 P 侧时触发,核心逻辑:

  1. 等待 P 端 block 就绪send_meta.ready.wait()——P 侧的 forward pass 可能尚未完成
  2. 地址计算_build_transfer_params 第 14241614 行):
    • 逐层、逐 block 计算 P 侧源地址和 D 侧目标地址
    • 处理 TP 异构(不同 TP 大小之间映射)
    • 处理 HMAHybrid Memory Allocator)多 KV group
    • 处理 Mamba/GDN 状态块
  3. RDMA 传输_send_blocks 第 16221651 行):
    ret_value = self.engine.batch_transfer_sync_write(
        remote_session, src_ptrs, dst_ptrs, lengths
    )  
    
    P 侧直接通过 RDMA/NVLink 将 KV Cache 字节写入 D 侧的 GPU 内存。
  4. 传输完成后,标记 finished_sending_reqs,释放 P 侧的 block。

五、KV Cache 传输完成后:请求恢复执行

5.1 Scheduler 收到传输完成信号

update_from_output_update_from_kv_xfer_finished(第 27132740 行):

for req_id in kv_connector_output.finished_recving or ():  
    req = self.requests[req_id]  
    if req.status == RequestStatus.WAITING_FOR_REMOTE_KVS:  
        self.finished_recving_kv_req_ids.add(req_id)  

5.2 下一轮调度时提升请求状态

_try_promote_blocked_waiting_request(第 26772711 行):

if request.status == RequestStatus.WAITING_FOR_REMOTE_KVS:  
    if request.request_id not in self.finished_recving_kv_req_ids:  
        return False  
    self._update_waiting_for_remote_kv(request)  
    request.status = RequestStatus.WAITING  # 回到正常调度队列  

_update_waiting_for_remote_kv(第 26342675 行):

  • 将接收到的 KV block 写入 prefix cachekv_cache_manager.cache_blocks
  • 如果全部 prompt token 都命中,减少 1 个 computed token 以确保最后 1 个 token 被本地重算(用于生成第一个 output token

5.3 正常调度执行

此时 request 回到 WAITING 状态,在下一轮 schedule() 中被正常调度:

  • num_computed_tokens 已经包含了外部加载的 token 数
  • Scheduler 从 num_computed_tokens 之后开始分配新的 compute
  • 新分配的 token(通常只有 1 个 decode token)在本地 GPU 上执行

六、完整时序总结

┌─ Proxy ──────────────────────────────────────────────────────────────┐  
│  ① POST /v1/chat/completions {"messages": [...], "max_tokens": 100}  │  
│  ② asyncio.create_task()        ③ Stream Decode Response             │  
│     → P Node (fire & forget)      → D Node (等待返回)                 │  
│     kv_transfer_params:           kv_transfer_params:                 │  
│       do_remote_decode=True         do_remote_prefill=True            │  
│       max_tokens=1                  remote_bootstrap_addr=<P地址>      │  
│       stream=False                  remote_engine_id=<P engine id>     │  
│       transfer_id=xfer-UUID         transfer_id=xfer-UUID             │  
└───────┬─────────────────────────────────┬───────────────────────────┘  
┌───────▼── P Node ───────────┐  ┌───────▼── D Node ────────────────────┐  
│ ④ Request 进入 Scheduler     │  │ ⑤ Request 进入 Scheduler             │  
│   is_kv_producer=True        │  │   is_kv_consumer=True                │  
│ ⑥ get_num_new_matched_tokens │  │ ⑦ get_num_new_matched_tokens        │  
│   → (0, False) 正常计算      │  │   → (N_prompt, True) 远端加载        │  
│ ⑧ Scheduler 分配 block       │  │ ⑨ Scheduler 分配 block              │  
│   → RUNNING 状态             │  │   → WAITING_FOR_REMOTE_KVS 状态     │  
│ ⑩ GPU forward pass           │  │ ⑪ Worker.start_load_kv()            │  
│   KV Cache 写入本地显存      │  │   → ZMQ 发送 MooncakeXferMetadata    │  
│ ⑫ send_kv_to_decode()        │  │   (含 D 侧 block 物理地址)          │  
│   ← ZMQ 收到 D 的请求        │◄───── ZMQ ────────────────────────────│  
│   → batch_transfer_sync_write│  │                                      │  
│   → RDMA 写入 D 侧 GPU 显存  ═══════ RDMA/NVLink ═══════════════════► │  
│ ⑬ finished_sending_reqs      │  │ ⑭ finished_recving_reqs             │  
│   释放 P 侧 block            │  │ ⑮ WAITING_FOR_REMOTE_KVS → WAITING  │  
│                             │  │ ⑯ 下一轮 schedule() 正常调度 decode  │  
│                             │  │ ⑰ GPU forward pass 自回归生成 output │  
│                             │  │ ⑱ 流式返回 tokens → Proxy → Client   │  
└──────────────────────────────┘  └──────────────────────────────────────┘  

七、关键设计要点

  1. P 和 D 收到的 prompt 完全一致,但通过 kv_transfer_params 标志区分行为。P 正常计算整个 prompt,D 跳过本地计算而从 P 拉取 KV Cache。
  2. transfer_id 是连接 P/D 的关键Proxy 生成 UUID,同时发给 P 和 D。P 侧通过 reqs_need_send[transfer_id] 等待 D 侧通过 ZMQ 发来的拉取请求(携带同样的 transfer_id),从而匹配同一请求。
  3. D 侧 block 物理地址是传输的基础D 侧在 receive_kv_from_single_worker 中将自己的 kv_caches_base_addr + block_lens + block_ids 发给 P,P 侧据此计算出 RDMA 的目标地址,直接写入 D 的 GPU 显存。
  4. 异步非阻塞设计D 侧在等待 KV 传输时处于 WAITING_FOR_REMOTE_KVS 状态,不占用 GPU 计算资源。同时 P 侧的发送也通过独立的后台线程池异步执行,不阻塞 P 的 forward pass。
  5. remote_engine_id 区分不同 DP rank 的 KV Cache:在 DP(数据并行)场景下,每个 DP rank 有独立的 KV Cache 分区。remote_engine_id 包含 engine_id_dp{N} 后缀,精确定位到哪个 DP rank 的 KV Cache。