跳到主要内容

数据截至 (上游 commit f25c580af159)

05 · 分布式执行与 API 服务

这一章讲什么: 同一份引擎代码怎么从单卡跑到上百卡,又怎么被包装成一个 OpenAI 兼容的 HTTP 服务。两条线:向下的 Executor(核心 → Worker 集群)和向上的 entrypoints(用户 → 前端)。


1. 它要解决的小问题

一个模型常常一张卡装不下(或装了但太慢),于是要拆:

并行方式拆的是什么代价
张量并行 TP每层的权重矩阵横向切开每层两次 all-reduce,通信密集,通常限单机
流水并行 PP层按段分给不同卡气泡;中间激活要跨卡传
数据并行 DP整个模型复制多份,各吃一部分请求显存 ×N,但互不干扰
专家并行 EPMoE 的专家分到不同卡all-to-all 通信

这些方案对上层应该是透明的:调度器和 EngineCore 不该知道模型拆成了几份。同时,对外服务还要接 HTTP、鉴权、SSE 流式、OpenAI 协议——这些也不该污染引擎。


2. 思路:一个方向一层壳

  • 向下:Executor 抽象。EngineCore 只调 execute_model(scheduler_output),背后是 1 个进程还是 32 个进程,Executor 内部消化。
  • 向上:AsyncLLM 已经够用(异步、流式),HTTP 层只是把请求翻成 AsyncLLM.generate(...) 调用、把输出翻成 SSE。

图示:一次 execute_model 的多进程展开

EngineCore WorkerProc × N(每卡一个)
────────── ────────────────────────
execute_model(scheduler_output)
│ collective_rpc:
▼ 方法名+参数 入队
rpc_broadcast_mq(共享内存)
│ ┌─ while True:
│ ─────────────────────────────► │ mq.dequeue()
│ │ getattr(worker, method)(*args)
│ │ → Worker.execute_model
│ │ → GPUModelRunner.execute_model
│ └─ 仅 output_rank 回传 ModelRunnerOutput

Future.result() 拿到聚合结果

怎么读这张图: 左到右是一次广播。注意不是每个 Worker 都回话——forward 所需的 all-reduce 在 Worker 之间用 NCCL 完成,回到 EngineCore 的只有指定 rank 的一份结果。


3. 原理演示:collective_rpc 的形状

# 示意,非源码
class MultiprocExecutor:
def execute_model(self, scheduler_output, non_block=False):
return self.collective_rpc(
"execute_model", # 方法名,worker 侧 getattr
args=(scheduler_output,),
unique_reply_rank=self.output_rank, # 只等这一个 rank 的回包
non_block=non_block, # True → 返回 Future
)

class WorkerProc:
def worker_busy_loop(self):
while True:
method, args, kwargs, out_rank = self.rpc_mq.dequeue()
result = getattr(self.worker, method)(*args, **kwargs)
if out_rank is None or self.rank == out_rank:
self.send_output(result)

重点:调度器输出只序列化一次、广播给全体;Worker 是纯执行体,没有任何调度决策权。


4. 真实实现

4.1 Executor 选型

Executor.get_classvllm/v1/executor/abstract.py:49)按配置挑实现:

实现场景位置
UniProcExecutor单卡、同进程(零 IPC 开销)vllm/v1/executor/uniproc_executor.py:51
MultiprocExecutor多卡单机默认:每卡一个子进程vllm/v1/executor/multiproc_executor.py:111
RayDistributedExecutor / RayExecutorV2Ray 集群、多机vllm/v1/executor/ray_executor.py:67ray_executor_v2.py:237
ExecutorWithExternalLauncher外部编排(如 MPI/torchrun)vllm/v1/executor/uniproc_executor.py:161

4.2 MultiprocExecutor 的通信骨架

  • 下行(核心→Worker):rpc_broadcast_mqvllm/v1/executor/multiproc_executor.py:144),共享内存消息队列;collective_rpc(:375)把 (方法名, args, kwargs, output_rank) 入队(:420)。execute_model(:340)就是它的特化。
  • Worker 侧WorkerProc.worker_busy_loop(:1029)死循环 dequeue → 执行 → 仅 output_rank 回传(:1049-1051)。每个 Worker 持有 GPUWorkervllm/v1/worker/gpu_worker.py),再往下是上一章的 GPUModelRunner
  • 失败传播:Worker 异常被包成 FAILURE 响应回传,executor 置 is_failed,后续 RPC 直接 RuntimeError("Executor failed.")(:392)。

TP/PP 的组初始化在 Worker 进程内完成(init_distributed_environment + 模型加载时按 rank 切权重);NCCL 集合通信对 EngineCore 完全透明。DP 则是另一种摆法:每个 DP rank 是一个完整的 EngineCore(各带自己的调度器和 KV 池),前端用 DPLBAsyncMPClient 路由(第三章)。

4.3 向上:vllm serve 的包装链

调用链(每环一层薄壳):

干什么位置
CLI解析参数,ServeSubcommandvllm/entrypoints/cli/serve.py:45
启动器建 socket、起 uvicorn、run_servervllm/entrypoints/launchers/api_server/entry.py:163
FastAPI app路由注册、app statevllm/entrypoints/launchers/app.py:19
协议层OpenAI chat/completion 协议解析、tool call、reasoning 拆分vllm/entrypoints/openai/chat_completion/serving.py:116OpenAIServingChat
引擎前端AsyncLLM(见第三章)vllm/v1/engine/async_llm.py:72

流式输出的通路:协议层拿到 AsyncLLM.generate 的异步迭代器 → 每个 RequestOutput 增量翻成 SSE data: 帧。旧的 python -m vllm.entrypoints.openai.api_server 入口在本 commit 已标 deprecated,文件顶部就是 DeprecationWarning,改用 vllm servevllm/entrypoints/openai/api_server.py:1-59)。

4.4 离线形态

不走 HTTP 时,LLM 类(vllm/entrypoints/llm.py:67)直接内嵌 LLMEngine +(可选)同进程 InprocClient,批量 generate 同步返回。和在线服务共用同一引擎核心与调度器,只是外皮不同。


5. 关键细节与坑

  • TP 对通信延迟敏感:每层 all-reduce 两次,跨机 TP 通常不划算;多机部署的惯用组合是机内 TP × 机间 PP 或 DP。
  • unique_reply_rank 是省流量的关键:N 卡 TP 下只有 1 个 rank 回传结果,否则 EngineCore 每步要收 N 份重复输出。
  • Worker 是哑终端:所有智能(调度、KV 分配、停止条件)都在 EngineCore;往 Worker 加状态前要三思——第三章的 async scheduling 要求 Worker 侧不假设「上一步结果已知」。
  • 多机靠 Ray 或外部编排,vLLM 自己不实现集群管理;external_launcher 模式把进程编排完全交给用户。
  • 协议层厚而引擎薄:tool calling、reasoning 解析、logprobs 格式都在 vllm/entrypoints/ 一侧完成,引擎只出 token。改协议行为别去动引擎。