Chengshu@skadai · 2026.09.30
2,004 字 · 2,954 词 · 约 13 分钟

走进 vLLM(四):从单卡到多机——MultiProcExecutor 与分布式服务

权重装不下一张卡之后怎么办:MultiProcExecutor 怎么用共享内存消息队列把 8 个 worker 进程编排起来;DP 复制、DPCoordinator 与 DP wave 又怎么把两个 8×H100 节点拼成一套服务;最后完整走一遍一条 curl 请求的生命周期。

本篇属于系列 走进 vLLM:高吞吐 LLM 推理系统解剖 · 第 4 篇

原文 · English中文译文

原文:Inside vLLM: Anatomy of a High-Throughput LLM Inference System,作者 Aleksa Gordic,2025-09-05 发布于 vLLM 官方博客。

本文是系列《走进 vLLM:高吞吐 LLM 推理系统解剖》的第 4 篇,共 5 篇。承接第 3 篇(speculative decoding 与 disaggregated P/D)。

本站中英对照排版:左栏英文原文,右栏中文译文;不好翻译的术语保留英文写法。

From UniprocExecutor to MultiProcExecutor

From UniprocExecutor to MultiProcExecutor

With the core techniques in place, we can now talk about scaling up.

核心技巧都就位了,现在可以聊横向扩容(scale up)了。

Suppose your model weights no longer fit into a single GPU’s VRAM.

假设你的模型权重已经装不进单张 GPU 的显存。

The first option is to shard the model across multiple GPUs on the same node using tensor parallelism (e.g., TP=8). If the model still doesn’t fit, the next step is pipeline parallelism across nodes.

第一个选项是用 tensor parallelism 把模型切到同节点内的多张卡上(比如 TP=8)。如果还是装不下,下一步就是跨节点的 pipeline parallelism。

Notes:

  • Intranode bandwidth is significantly higher than internode, which is why tensor parallelism (TP) is generally preferred over pipeline parallelism (PP). (It is also true that PP communicates less data than TP.)
  • I’m not covering expert parallelism (EP) since we’re focusing on standard transformers rather than MoE, nor sequence parallelism, as TP and PP are the most commonly used in practice.

说明

  • 节点内带宽显著高于节点间带宽,所以 tensor parallelism(TP)通常优于 pipeline parallelism(PP)。(同样成立的是:PP 传输的数据量比 TP 少。)
  • 这里不覆盖 expert parallelism(EP),因为我们聚焦的是标准 transformer 而不是 MoE;也不覆盖 sequence parallelism,因为实践中最常用的就是 TP 和 PP。

At this stage, we need multiple GPU processes (workers) and an orchestration layer to coordinate them. That’s exactly what MultiProcExecutor provides.

到了这个阶段,我们需要多个 GPU 进程(worker),以及一个编排层来协调它们。这正是 MultiProcExecutor 提供的。

图 13:TP=8 场景下的 MultiProcExecutor(driver worker 是 rank 0)
图 13:TP=8 场景下的 MultiProcExecutor(driver worker 是 rank 0)

Figure 13: MultiProcExecutor in a TP=8 setting (driver worker being rank 0)

How this works in vLLM:

在 vLLM 里这是怎么实现的:

  1. MultiProcExecutor initializes an rpc_broadcast_mq message queue (implemented with shared memory under the hood).
  2. The constructor loops over world_size (e.g. TP=8 ⇒ world_size=8) and spawns a daemon process for each rank via WorkerProc.make_worker_process.
  3. For each worker, the parent first creates a reader and writer pipe.
  4. The new process runs WorkerProc.worker_main, which instantiates a worker (going through the same “init device”, “load model”, etc. as in UniprocExecutor).
  5. Each worker determines whether it is the driver (rank 0 in the TP group) or a regular worker. Every worker sets up two queues:
  1. MultiProcExecutor 初始化一个 rpc_broadcast_mq 消息队列(底层用共享内存实现)。
  2. 构造函数遍历 world_size(比如 TP=8 ⇒ world_size=8),为每个 rank 通过 WorkerProc.make_worker_process 起一个 daemon 进程。
  3. 对每个 worker,父进程先创建一对 reader / writer pipe。
  4. 新进程运行 WorkerProc.worker_main,它实例化一个 worker(走的是和 UniprocExecutor 里一样的 “init device”、“load model” 等流程)。
  5. 每个 worker 判断自己是 driver(TP 组里的 rank 0)还是普通 worker。每个 worker 都建两个队列:
  • rpc_broadcast_mq (shared with the parent) for receiving work.
  • worker_response_mq for sending responses back.
  • rpc_broadcast_mq(与父进程共享),用来接收工作。
  • worker_response_mq,用来把响应发回去。
  1. During initialization, each child sends its worker_response_mq handle to the parent via the pipe. Once all are received, the parent unblocks — this completes coordination.
  2. Workers then enter a busy loop, blocking on rpc_broadcast_mq.dequeue. When a work item arrives, they execute it (just like in UniprocExecutor, but now with TP/PP-specific partitioned work). Results are sent back through worker_response_mq.enqueue.
  3. At runtime, when a request arrives, MultiProcExecutor enqueues it into rpc_broadcast_mq (non-blocking) for all children workers. It then waits on the designated output rank’s worker_response_mq.dequeue to collect the final result.
  1. 初始化期间,每个子进程通过 pipe 把自己的 worker_response_mq handle 发给父进程。全部收齐后父进程解除阻塞——协调完成。
  2. 之后 worker 进入忙循环,阻塞在 rpc_broadcast_mq.dequeue 上。有工作项到达就执行它(和 UniprocExecutor 里一样,只不过现在的工作按 TP/PP 做了切分)。结果通过 worker_response_mq.enqueue 发回。
  3. 运行时,请求到达后 MultiProcExecutor 把它(非阻塞地)塞进 rpc_broadcast_mq,发给所有子 worker;然后在指定的 output rank 的 worker_response_mq.dequeue 上等待,收集最终结果。

From the engine’s perspective, nothing has changed — all of this multiprocessing complexity is abstracted away through a call to model executor’s execute_model.

从 engine 的视角看,什么都没变——所有这些多进程的复杂度都被抽象掉了,对外只是一个 model executor 的 execute_model 调用。

  • In the UniProcExecutor case: execute_model directly leads to calling execute_model on the worker
  • In the MultiProcExecutor case: execute_model indirectly leads to calling execute_model on each worker through rpc_broadcast_mq
  • 在 UniProcExecutor 的情况下:execute_model 直接导致 worker 上的 execute_model 被调用
  • 在 MultiProcExecutor 的情况下:execute_model 间接地通过 rpc_broadcast_mq 导致每个 worker 上的 execute_model 被调用

At this point, we can run models that are as large as resources allow using the same engine interface.

到这里,我们能用同一套 engine 接口,跑资源允许范围内任意大的模型。

The next step is to scale out: enable data parallelism (DP > 1) replicating the model across nodes, add a lightweight DP coordination layer, introduce load balancing across replicas, and place one or more API servers in front to handle incoming traffic.

下一步是横向扩展(scale out):启用 data parallelism(DP > 1)把模型复制到多个节点,加上一个轻量的 DP 协调层,引入副本之间的 load balancing,并在前面放一到多个 API server 来接流量。

Distributed system serving vLLM

Distributed system serving vLLM

There are many ways to set up serving infrastructure, but to stay concrete, here’s one example: suppose we have two H100 nodes and want to run four vLLM engines across them.

部署服务基础设施的方式有很多,为了讲得具体,这里给一个例子:假设我们有两个 H100 节点,想在它们上面跑四个 vLLM engine。

If the model requires TP=4, we can configure the nodes like this.

如果模型需要 TP=4,我们可以这样配置这两个节点。

图 14:2 个 8×H100 节点的服务配置(1 个 headless,1 个 API server)
图 14:2 个 8×H100 节点的服务配置(1 个 headless,1 个 API server)

Figure 14: server configuration with 2 8xH100 nodes (1 headless, 1 api server)

On the first node, run the engine in headless mode (no API server) with the following arguments:

在第一个节点上,用 headless 模式(不带 API server)启动 engine,参数如下:

vllm serve <model-name>
  --tensor-parallel-size 4
  --data-parallel-size 4
  --data-parallel-size-local 2
  --data-parallel-start-rank 0
  --data-parallel-address <master-ip>
  --data-parallel-rpc-port 13345
  --headless

and run that same command on the other node with few tweaks:

在另一个节点上跑同一条命令,只改两处:

  • no –headless
  • modify DP start rank
  • 不带 --headless
  • 修改 DP start rank
vllm serve <model-name>
  --tensor-parallel-size 4
  --data-parallel-size 4
  --data-parallel-size-local 2
  --data-parallel-start-rank 2
  --data-parallel-address <master-ip>
  --data-parallel-rpc-port 13345

说明

This assumes networking is configured so all nodes can reach the specified IP and port.

说明

这里假设网络已经配好,所有节点都能访问指定的 IP 和端口。

How does this work in VLLM?

在 vLLM 里这是怎么实现的?

On the headless server node

On the headless server node

On the headless node, a CoreEngineProcManager launches 2 processes (per –data-parallel-size-local) each running EngineCoreProc.run_engine_core. Each of these functions creates a DPEngineCoreProc (the engine core) and then enters its busy loop.

在 headless 节点上,CoreEngineProcManager 会起 2 个进程(数量由 --data-parallel-size-local 决定),每个进程运行 EngineCoreProc.run_engine_core。这些函数各自创建一个 DPEngineCoreProc(engine core),然后进入它的忙循环。

DPEngineCoreProc initializes its parent EngineCoreProc (child of EngineCore), which:

DPEngineCoreProc 初始化它的父类 EngineCoreProc(EngineCore 的子类),后者会:

  1. Creates an input_queue and output_queue (queue.Queue).
  2. Performs an initial handshake with the frontend on the other node using a DEALER ZMQ socket (async messaging lib), and receives coordination address info.
  3. Initializes DP group (e.g. using NCCL backend).
  4. Initializes the EngineCore with MultiProcExecutor (TP=4 on 4 GPUs as described earlier).
  5. Creates a ready_event (threading.Event).
  6. Starts an input deamon thread (threading.Thread) running process_input_sockets(…, ready_event). Similarly starts an output thread.
  7. Still in the main thread, waits on ready_event until all input threads across all 4 processes (spanning the 2 nodes) have completed the coordination handshake finally executing ready_event.set().
  8. Once unblocked, sends a “ready” message to the frontend with metadata (e.g., num_gpu_blocks available in paged KV cache memory).
  9. The main, input, and output threads then enter their respective busy loops.
  1. 创建 input_queue 和 output_queue(queue.Queue)。
  2. 用 DEALER ZMQ socket(异步消息库)与另一个节点上的 frontend 做一次初始握手,并接收协调地址信息。
  3. 初始化 DP group(比如用 NCCL 后端)。
  4. 用 MultiProcExecutor 初始化 EngineCore(就是前面说的 4 张卡上 TP=4)。
  5. 创建 ready_event(threading.Event)。
  6. 启动一个 input daemon 线程(threading.Thread)运行 process_input_sockets(…, ready_event);同样地启动一个 output 线程。
  7. 仍在主线程里等待 ready_event,直到所有 4 个进程(横跨 2 个节点)的 input 线程都完成协调握手,最后执行 ready_event.set()。
  8. 解除阻塞后,给 frontend 发一条带 metadata 的 “ready” 消息(比如 paged KV cache 显存里可用的 num_gpu_blocks)。
  9. 之后 main、input、output 三个线程各自进入忙循环。

TL;DR: We end up with 4 child processes (one per DP replica), each running a main, input, and output thread. They complete a coordination handshake with the DP coordinator and frontend, then all three threads per process run in steady-state busy loops.

TL;DR:我们最终得到 4 个子进程(每个 DP 副本一个),每个进程里跑着 main、input、output 三个线程。它们与 DP coordinator 和 frontend 完成一次协调握手,然后每个进程的三个线程都在稳态忙循环里跑着。

图 15:4 个 DP 副本、4 个 DPEngineCoreProc 的分布式系统
图 15:4 个 DP 副本、4 个 DPEngineCoreProc 的分布式系统

Figure 15: distributed system with 4 DP replicas running 4 DPEngineCoreProc

Current steady state:

当前稳态:

  • Input thread — blocks on the input socket until a request is routed from the API server; upon receipt, it decodes the payload, enqueues a work item via input_queue.put_nowait(…), and returns to blocking on the socket.
  • Main thread — wakes on input_queue.get(…), feeds the request to the engine; MultiProcExecutor runs the forward pass and enqueues results to output_queue.
  • Output thread — wakes on output_queue.get(…), sends the result back to the API server, then resumes blocking.
  • Input thread —— 阻塞在 input socket 上,直到有请求从 API server 路由过来;收到后解码 payload,通过 input_queue.put_nowait(...) 入队一个工作项,然后回到 socket 上继续阻塞。
  • Main thread —— 在 input_queue.get(...) 上被唤醒,把请求喂给 engine;MultiProcExecutor 跑 forward pass,并把结果入队到 output_queue。
  • Output thread —— 在 output_queue.get(...) 上被唤醒,把结果发回 API server,然后继续阻塞。

Additional mechanics:

其余机制:

  • DP wave counter — the system tracks “waves”; when all engines become idle they quiesce, and the counter increments when new work arrives (useful for coordination/metrics).
  • Control messages — the API server can send more than just inference requests (e.g., aborts and utility/control RPCs).
  • Dummy steps for lockstep — if any DP replica has work, all replicas execute a forward step; replicas without requests perform a dummy step to participate in required synchronization points (avoids blocking the active replica).
  • DP wave counter —— 系统会跟踪 “wave”;当所有 engine 都空闲下来时它们静默,新工作到达时计数器递增(对协调/指标有用)。
  • Control messages —— API server 能发的不只是推理请求(比如 abort,以及工具/控制类 RPC)。
  • Dummy steps for lockstep —— 只要有任何一个 DP 副本有工作,所有副本都要执行一次 forward step;没有请求的副本会跑一次 dummy step,以参与必需的同步点(避免把活跃副本卡住)。

说明

Lockstep clarification: this is actually only required for MoE models where the expert layers form an EP or TP group while attention layers are still DP. It’s currently always done with DP - this is just because there’s limited use for “built-in” non-MoE DP since you could just run multiple independent vLLMs and load-balance between them in a normal way.

说明

关于 lockstep 的澄清:实际上只有 MoE 模型才需要——它的 expert 层构成一个 EP 或 TP 组,而 attention 层仍是 DP。目前它对所有 DP 都会这么做,只是因为“内置的非 MoE DP”用处有限:你完全可以直接跑多个独立的 vLLM,用常规方式在它们之间做负载均衡。

Now for the second part, what happens on the API server node?

接下来是第二部分:API server 节点上发生了什么?

On the API server node

On the API server node

We instantiate an AsyncLLM object (an asyncio wrapper around the LLM engine). Internally this creates a DPLBAsyncMPClient (data-parallel, load-balancing, asynchronous, multiprocessing client).

我们实例化一个 AsyncLLM 对象(LLM engine 的 asyncio 封装)。它内部会创建一个 DPLBAsyncMPClient(data-parallel、load-balancing、asynchronous、multiprocessing client)。

Inside the parent class of MPClient, the launch_core_engines function runs and:

在 MPClient 的父类内部,launch_core_engines 函数会运行并:

  1. Creates the ZMQ addresses used for the startup handshake (as seen on the headless node).
  2. Spawns a DPCoordinator process.
  3. Creates a CoreEngineProcManager (same as on the headless node).
  1. 创建启动握手用的 ZMQ 地址(和 headless 节点上看到的一样)。
  2. 起一个 DPCoordinator 进程。
  3. 创建一个 CoreEngineProcManager(和 headless 节点上相同)。

Inside AsyncMPClient (child of MPClient), we:

在 AsyncMPClient(MPClient 的子类)里,我们:

  1. Create an outputs_queue (asyncio.Queue).
  2. We create an asyncio task process_outputs_socket which communicates (through the output socket) with output threads of all 4 DPEngineCoreProc and writes into outputs_queue.
  3. Subsequently one more asyncio task output_handler from AsyncLLM reads from this queue and finally sends out information to the create_completion function.
  1. 创建 outputs_queue(asyncio.Queue)。
  2. 创建一个 asyncio task process_outputs_socket,它通过 output socket 与全部 4 个 DPEngineCoreProc 的 output 线程通信,并写进 outputs_queue。
  3. 随后,来自 AsyncLLM 的另一个 asyncio task output_handler 从这个队列读取,最终把信息送进 create_completion 函数。

Inside DPAsyncMPClient we create an asyncio task run_engine_stats_update_task which communicates with DP coordinator.

在 DPAsyncMPClient 里,我们创建一个 asyncio task run_engine_stats_update_task,它与 DP coordinator 通信。

The DP coordinator mediates between the frontend (API server) and backend (engine cores). It:

DP coordinator 在 frontend(API server)和 backend(engine core)之间做中介。它会:

  • Periodically sends load-balancing info (queue sizes, waiting/running requests) to the frontend’s run_engine_stats_update_task.
  • Handles SCALE_ELASTIC_EP commands from the frontend by dynamically changing the number of engines (only works with Ray backend).
  • Sends START_DP_WAVE events to the backend (when triggered by frontend) and reports wave-state updates back.
  • 定期把负载均衡信息(队列长度、waiting/running 请求数)发给 frontend 的 run_engine_stats_update_task。
  • 处理来自 frontend 的 SCALE_ELASTIC_EP 命令,动态改变 engine 数量(只在 Ray 后端下可用)。
  • 向 backend 发送 START_DP_WAVE 事件(由 frontend 触发),并回报 wave 状态更新。

To recap, the frontend (AsyncLLM) runs several asyncio tasks (remember: concurrent, not parallel):

小结一下,frontend(AsyncLLM)跑着若干个 asyncio task(记住:是并发,不是并行):

  • A class of tasks handles input requests through the generate path (each new client request spawns a new asyncio task).
  • Two tasks (process_outputs_socket, output_handler) process output messages from the underlying engines.
  • One task (run_engine_stats_update_task) maintains communication with the DP coordinator: sending wave triggers, polling LB state, and handling dynamic scaling requests.
  • 一类 task 通过 generate 路径处理输入请求(每个新的客户端请求都派生一个新的 asyncio task)。
  • 两个 task(process_outputs_socket、output_handler)处理来自底层 engine 的输出消息。
  • 一个 task(run_engine_stats_update_task)维持与 DP coordinator 的通信:发送 wave 触发、轮询 LB 状态、处理动态扩容请求。

Finally, the main server process creates a FastAPI app and mounts endpoints such as OpenAIServingCompletion and OpenAIServingChat, which expose /completion, /chat/completion, and others. The stack is then served via Uvicorn.

最后,主 server 进程创建一个 FastAPI app,挂上 OpenAIServingCompletion、OpenAIServingChat 这类 endpoint,对外暴露 /completion、/chat/completion 等接口。整个栈由 Uvicorn 提供服务。

So, putting it all together, here’s the full request lifecycle!

把这些拼起来,就是一条请求的完整生命周期!

You send from your terminal:

你在终端里发出:

curl -X POST http://localhost:8000/v1/completions -H "Content-Type: application/json" -d '{
  "model": "TinyLlama/TinyLlama-1.1B-Chat-v1.0",
  "prompt": "The capital of France is",
  "max_tokens": 50,
  "temperature": 0.7
}'

What happens next:

接着会发生:

  1. The request hits OpenAIServingCompletion’s create_completion route on the API server.
  2. The function tokenizes the prompt asynchronously, and prepares metadata (request ID, sampling params, timestamp, etc.).
  3. It then calls AsyncLLM.generate, which follows the same flow as the synchronous engine, eventually invoking DPAsyncMPClient.add_request_async.
  4. This in turn calls get_core_engine_for_request, which does load balancing across engines based on the DP coordinator’s state (picking the one that has minimal score / lowest load: score = len(waiting) * 4 + len(running)).
  5. The ADD request is sent to the chosen engine’s input_socket.
  6. At that engine:
  1. 请求命中 API server 上 OpenAIServingCompletion 的 create_completion 路由。
  2. 该函数异步地对 prompt 做 tokenization,并准备 metadata(request ID、sampling params、时间戳等)。
  3. 然后它调用 AsyncLLM.generate,之后的流程和同步 engine 一样,最终调用 DPAsyncMPClient.add_request_async。
  4. 这又调用 get_core_engine_for_request,根据 DP coordinator 的状态在各个 engine 之间做负载均衡(挑分数最小、也就是负载最低的那个:score = len(waiting) * 4 + len(running))。
  5. ADD 请求被发到所选 engine 的 input_socket。
  6. 在那个 engine 上:
  • Input thread — unblocks, decodes data from the input socket, and places a work item on the input_queue for the main thread.
  • Main thread — unblocks on input_queue, adds the request to the engine, and repeatedly calls engine_core.step(), enqueueing intermediate results to output_queue until a stop condition is met.
  • Input thread —— 解除阻塞,从 input socket 解码数据,把一个工作项放进 input_queue 给主线程。
  • Main thread —— 在 input_queue 上解除阻塞,把请求加入 engine,并反复调用 engine_core.step(),把中间结果入队到 output_queue,直到满足停止条件。

说明

Reminder: step() calls the scheduler, model executor (which in turn can be MultiProcExecutor!), etc. We have already seen this!

说明

提醒一下:step() 会调用 scheduler、model executor(后者又可能是 MultiProcExecutor!)等等。这些我们前面都见过了!

  • Output thread — unblocks on output_queue and sends results back through the output socket.
  • Output thread —— 在 output_queue 上解除阻塞,把结果通过 output socket 发回。
  1. Those results trigger the AsyncLLM output asyncio tasks (process_outputs_socket and output_handler), which propagate tokens back to FastAPI’s create_completion route.
  2. FastAPI attaches metadata (finish reason, logprobs, usage info, etc.) and returns a JSONResponse via Uvicorn to your terminal!
  1. 这些结果触发 AsyncLLM 的 output asyncio task(process_outputs_socket 和 output_handler),它们把 token 一路送回 FastAPI 的 create_completion 路由。
  2. FastAPI 附上 metadata(finish reason、logprobs、usage 信息等),通过 Uvicorn 返回一个 JSONResponse 到你终端!

And just like that, your completion came back — the whole distributed machinery hidden behind a simple curl command! :) So much fun!!!

就这样,你的 completion 回来了——整套分布式机器都藏在一条简单的 curl 命令后面!:) 太有趣了!!!

Additional notes:

  • When adding more API servers, load balancing is handled at the OS/socket level. From the application’s perspective, nothing significant changes — the complexity is hidden.
  • With Ray as a DP backend, you can expose a URL endpoint (/scale_elastic_ep) that enables automatic scaling of the number of engine replicas up or down.

补充说明

  • 增加更多 API server 时,负载均衡是在 OS/socket 层面处理的。从应用视角看没什么显著变化——复杂度被隐藏了。
  • 用 Ray 作为 DP 后端时,你可以暴露一个 URL endpoint(/scale_elastic_ep)来启用自动扩缩容。

讨论

这里是静态站点,没有内嵌评论区。如果这篇文章对你有用,欢迎通过 RSS 订阅后续更新。