- 46 篇原文+摘要双文件归档(按 来源/作者 分层,复用本地归档 20 篇+新抓取 26 篇) - 即梦生成 9 组主题配图(大图+列表缩略图)存入 知识/金鹏/20260806/ - 章节重组为 9 个主题分组并挂接摘要引用
13 KiB
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, 第 250–280 行):
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, 第 283–312 行):
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_id(UUID)——这是 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 第 116–119 行:
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.py 的 schedule() 方法中。
3.1 Prefill 节点 (is_kv_producer) 的调度
当 do_remote_decode=True 时:
get_num_new_matched_tokens(MooncakeConnectorScheduler 第 712–755 行):返回(0, False)——P 节点不需要加载任何外部 KV,正常计算 prompt。update_state_after_alloc(第 757–805 行):记录该请求需要将 KV Cache 发送出去:elif params.get("do_remote_decode"): self._reqs_need_send[request.request_id] = (request, [])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 时:
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 = 异步加载- Scheduler 分配 KV block(第 953–971 行):
allocate_slots中传入num_external_computed_tokens,并为远程 token 预留 block。关键参数:delay_cache_blocks=True— block 分配了但不立即写入 prefix cache(因为数据还在远端)reserved_blocks— 确保有足够空间完成整个异步传输,防止死锁
- 请求进入
WAITING_FOR_REMOTE_KVS状态(第 1010–1040 行):此时 request 被放入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 竞态skipped_waiting队列,不会进入running列表,本轮不参与 GPU 计算。 build_connector_meta:将需要 recv 的请求打包成MooncakeConnectorMetadata(含remote_engine_id、remote_bootstrap_addr、transfer_id、local_block_ids),放入SchedulerOutput.kv_connector_metadata传递给 Worker。
四、Worker 层:Mooncake 的 KV Cache RDMA 传输
4.1 D 节点 Worker:拉取 KV Cache
start_load_kv(MooncakeConnectorWorker, 第 2028–2039 行):
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_kv → handle_new_engine_id → receive_kv → receive_kv_from_single_worker:
- 首先通过
_connect_to_prefiller_bootstrap(remote_bootstrap_addr)查询 P 节点的 bootstrap server,获取 P 侧所有 TP/PP Worker 的 ZMQ 地址 - 构建
MooncakeXferMetadata(第 1825–1840 行),包含:remote_hostname/remote_port:D 侧的地址,P 侧用来做 RDMA 写入remote_tp_size/remote_tp_rank:TP 拓扑信息req_blocks:{req_id: (transfer_id, local_block_ids)}——D 侧已分配的 block IDkv_caches_base_addr/block_lens/kv_block_lens:D 侧 KV Cache 的物理内存地址和布局
- 通过 ZMQ DEALER socket 向 P 侧对应 TP rank 的 Worker 发送此元数据
- P 侧收到后,通过 RDMA (
batch_transfer_sync_write) 将 KV Cache 数据直接写入 D 侧的 GPU 内存地址 - 传输完成后,D 侧标记
finished_recving_reqs
4.2 P 节点 Worker:推送 KV Cache
record_send_reqs(第 2000–2026 行):P 侧的 Worker 将 Scheduler 传递来的 reqs_to_send(含 transfer_id 和 local_block_ids)记录到 self.reqs_need_send 字典。
send_kv_to_decode(第 1193–1382 行):当 D 侧的 ZMQ 请求到达 P 侧时触发,核心逻辑:
- 等待 P 端 block 就绪:
send_meta.ready.wait()——P 侧的 forward pass 可能尚未完成 - 地址计算(
_build_transfer_params第 1424–1614 行):- 逐层、逐 block 计算 P 侧源地址和 D 侧目标地址
- 处理 TP 异构(不同 TP 大小之间映射)
- 处理 HMA(Hybrid Memory Allocator)多 KV group
- 处理 Mamba/GDN 状态块
- RDMA 传输(
_send_blocks第 1622–1651 行):P 侧直接通过 RDMA/NVLink 将 KV Cache 字节写入 D 侧的 GPU 内存。ret_value = self.engine.batch_transfer_sync_write( remote_session, src_ptrs, dst_ptrs, lengths ) - 传输完成后,标记
finished_sending_reqs,释放 P 侧的 block。
五、KV Cache 传输完成后:请求恢复执行
5.1 Scheduler 收到传输完成信号
update_from_output → _update_from_kv_xfer_finished(第 2713–2740 行):
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(第 2677–2711 行):
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(第 2634–2675 行):
- 将接收到的 KV block 写入 prefix cache(
kv_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 │
└──────────────────────────────┘ └──────────────────────────────────────┘
七、关键设计要点
- P 和 D 收到的 prompt 完全一致,但通过
kv_transfer_params标志区分行为。P 正常计算整个 prompt,D 跳过本地计算而从 P 拉取 KV Cache。 transfer_id是连接 P/D 的关键:Proxy 生成 UUID,同时发给 P 和 D。P 侧通过reqs_need_send[transfer_id]等待 D 侧通过 ZMQ 发来的拉取请求(携带同样的transfer_id),从而匹配同一请求。- D 侧 block 物理地址是传输的基础:D 侧在
receive_kv_from_single_worker中将自己的kv_caches_base_addr+block_lens+block_ids发给 P,P 侧据此计算出 RDMA 的目标地址,直接写入 D 的 GPU 显存。 - 异步非阻塞设计:D 侧在等待 KV 传输时处于
WAITING_FOR_REMOTE_KVS状态,不占用 GPU 计算资源。同时 P 侧的发送也通过独立的后台线程池异步执行,不阻塞 P 的 forward pass。 remote_engine_id区分不同 DP rank 的 KV Cache:在 DP(数据并行)场景下,每个 DP rank 有独立的 KV Cache 分区。remote_engine_id包含engine_id_dp{N}后缀,精确定位到哪个 DP rank 的 KV Cache。