EngineCoreClient:进程内与多进程两条客户端通路
职责
vLLM 的前端(LLM / AsyncLLM / serving 入口)不直接摸 EngineCore,而是经一层 client 把请求推进去、把输出拉出来。这一层就是 EngineCoreClient,它把「同步还是异步」「同进程还是子进程」这两条正交维度,统一到同一组抽象方法上。vllm/v1/engine/core_client.py 这个文件里塞下了所有实现。
最关键的两个实现是 InprocClient 和 MPClient。前者把 EngineCore 直接构造在本进程里,调一步算一步,没有 busy loop、没有 ZMQ,主要给 V0 风格的 LLMEngine.add_request()/step() 当薄壳。后者把 EngineCore 包成 EngineCoreProc 子进程,通过 ZMQ ROUTER/PULL 双 socket 收发请求与输出,前端调用线程只跟 socket 打交道。MPClient 是基类,具体落地有 SyncMPClient(同步,起后台线程从 output socket 收帧)和 AsyncMPClient(asyncio,起后台 task 跑同一个 socket),后者另页讲。
EngineCoreClient.make_client 是个静态工厂,根据 multiprocess_mode 和 asyncio_mode 两个布尔挑实现。组合 (asyncio=True, multiprocess=False) 直接抛 NotImplementedError——EngineCore 自己不是 asyncio-friendly 的,异步必须配多进程一起用。
设计动机
- 解耦前端与 EngineCore 的生命周期:同进程调 EngineCore 简单,但前端要并发跑多个 EngineCore(数据并行)、或者要做弹性 EP(在线加减 rank),就必须把 EngineCore 推到独立进程里,前端只持 socket。
MPClient.__init__里weakref.finalize(self, self.resources)这一句是核心:就算构造中途抛异常,后台子进程和 ZMQ context 也会被BackgroundResources.__call__兜底回收,不会泄漏。 - 同步接口和异步接口共享同一套 ZMQ 协议:
SyncMPClient用queue.Queue把 output socket 的帧搬到主线程,AsyncMPClient用asyncio.Queue搬,但MPClient共享了 encoder/decoder、ready 握手、engine monitor、utility future 注册表。差异被压在两个子类里。 - InprocClient 的
sleep(mode="wait")直接禁掉(InprocClient.sleepL324-L328),因为本进程模式下没有可以「等到空闲再 sleep」的并发调度器,只能abort。这种边界在子进程模式里是支持的——子进程有自己的 busy loop 和调度器状态机。 - utility 调用走 future 注册表:像
collective_rpc、save_sharded_state这类需要返回值的 RPC,不能像add_request那样 fire-and-forget。call_utility生成一个call_id,把 future 塞进self.utility_results,等UtilityOutput从 output socket 回来时_process_utility_output用call_id取出 future 并set_result。这套机制在 sync 和 async 两种 client 上完全一致,只是 future 类型不同(concurrent.futures.Futurevsasyncio.Future)。 - engine dead 传播:子进程意外死亡时,monitor 线程把
BackgroundResources.engine_dead置 True,后续任何_send_input都会被ensure_alive挡住抛EngineDeadError。output socket 那边如果收到ENGINE_CORE_DEAD单帧,validate_alive也会把同一标志位拉起来,保证发请求侧和收输出侧都看得到。
关键文件
EngineCoreClient + make_client:71-105— 抽象基类和静态工厂,根据multiprocess_mode/asyncio_mode选InprocClient/SyncMPClient/AsyncMPClientmake_async_mp_client:107-132— 异步工厂,按data_parallel_size与是否外部 LB 选AsyncMPClient/DPAsyncMPClient/DPLBAsyncMPClientInprocClient.__init__ / get_output:276-292— 直接构造EngineCore,get_output就是同步调step_fn()+post_step()InprocClient.add_request / abort / shutdown:297-306— 直接转发给self.engine_core,没有序列化、没有 socketInprocClient.sleep 拒绝 wait:324-328— 本进程模式不支持waitpause,只能abortBackgroundResources:370-458— dataclass +__call__作为weakref.finalize回调,统一关 socket / 停 task / shutdown engine managervalidate_alive:454-457— 收到单帧ENGINE_CORE_DEAD即标记engine_dead,抛EngineDeadErrorMPClient.__init__:480-649— 建 ZMQ context/socket、可选地launch_core_engines、握手 ready、起 monitor,失败时回滚MPClient.shutdown:651-662—detach()finalizer 后停 engine manager、跑BackgroundResources清理start_engine_core_monitor:685-712— daemon 线程monitor_engine_liveness,引擎死亡时置engine_dead并触发self.shutdown()_apply_ready_response:714-754— 解码EngineCoreReadyResponse,把max_model_len/num_gpu_blocks/block_size从子进程回填到前端vllm_configSyncMPClient.__init__:779-847— 建queue.Queue,起process_outputs_socketdaemon 线程轮询 output socket + shutdown socketSyncMPClient.get_output / _send_input / call_utility:849-881— 阻塞outputs_queue.get(),发请求时按需trackMessageTracker 防 tensor 提前回收
数据流
SyncMPClient 拉输出靠一个 daemon 线程在 ZMQ output_socket 上 poll,解码后塞进 queue.Queue。前端线程调 get_output() 就是 queue.get()——把 ZMQ 的非阻塞 socket 适配成 Python 同步阻塞接口。这段 process_outputs_socket 是整个同步通路的核心:
def process_outputs_socket():
assert isinstance(out_socket, zmq.Socket)
shutdown_socket = ctx.socket(zmq.PAIR)
try:
shutdown_socket.bind(shutdown_path)
poller = zmq.Poller()
poller.register(shutdown_socket, zmq.POLLIN)
poller.register(out_socket, zmq.POLLIN)
while True:
socks = poller.poll()
if not socks:
continue
if len(socks) == 2 or socks[0][0] == shutdown_socket:
# shutdown signal, exit thread.
break
frames = out_socket.recv_multipart(copy=False)
resources.validate_alive(frames)
outputs: EngineCoreOutputs = decoder.decode(frames)
if outputs.utility_output:
_process_utility_output(outputs.utility_output, utility_results)
else:
outputs_queue.put_nowait(outputs)
except Exception as e:
outputs_queue.put_nowait(e)
finally:
# Close sockets.
shutdown_socket.close(linger=0)
out_socket.close(linger=0)(SyncMPClient.process_outputs_socket:808-836)。注意三个细节:一是 shutdown 走单独的 zmq.PAIR inproc socket 绑 shutdown_path,BackgroundResources.__call__ 关闭时往这个 socket 发一个空字节通知线程退出,避免线程死等 out_socket。二是 validate_alive 把 ENGINE_CORE_DEAD 单帧转成 EngineDeadError。三是 utility 回包(如 collective_rpc 的返回)和普通 outputs 走同一帧,前者直接灌进 utility_results[call_id] 的 future,后者进队列。
发请求的方向更简单,SyncMPClient._send_input 直接把 (identity, request_type, *encoder.encode(request)) 多帧发到 ROUTER 上,有 tensor buffer 时用 track=True 拿到 MessageTracker,挂进 pending_messages 防止 Python 端把 backing tensor 提前 GC(_send_input:861-873)。
握手阶段在 MPClient.__init__ 里(等待 ready:615-634):对每个 engine_rank 生成 2 字节的 little-endian identity,在 input_socket 上 poll 等子进程的 EngineCoreReadyResponse,超时由 VLLM_ENGINE_READY_TIMEOUT_S 拉住。_apply_ready_response 把子进程跑出来的 num_gpu_blocks、对齐过的 block_size、自动拟合的 max_model_len 回填到前端的 vllm_config,这是前端真正能拿来调度和限流的真实容量。
边界与失败
asyncio_mode=True, multiprocess_mode=False直接拒绝 (make_client:91-95)。EngineCore 内部不是 async-friendly 的,这条组合没有实现。- 构造失败必须回收子进程:
MPClient.__init__用try/finally+success标志,失败时调self._finalizer()触发BackgroundResources.__call__,把已经launch_core_engines起来的子进程一起停掉(__init__ finally:499-649)。 - InprocClient 不支持
waitpause(InprocClient.sleep:324-328)。本进程模式下没有谁去等空闲,只能abort当前请求后 sleep。 - engine core 子进程死亡会被两个方向都感知:发送侧靠
ensure_alive(ensure_alive:670-672),接收侧靠validate_alive收到ENGINE_CORE_DEAD单帧(validate_alive:454-457)。两路都会把BackgroundResources.engine_dead拉起来。 - utility future 的 InvalidStateError 容忍:
_process_utility_output捕获asyncio.InvalidStateError(_process_utility_output:769-776),覆盖调用方 task 被取消导致 future 已 cancelled 的场景。 - output socket 关闭要在 context.term 之前:
BackgroundResources.__call__显式close_sockets再做后续(sync case:437-452),否则 ZMQ context 终止会卡住。 - SyncMPClient 的 output_socket 由线程负责关闭:构造最后把
self.resources.output_socket = None(socket ownership 转交:846-847),避免BackgroundResources和 daemon 线程同时关。 - ready 握手超时给可读报错:
Timed out waiting for engine core processes to start. This is often caused by slow weight loading for large models.(ready timeout:622-630)——直接告诉用户调VLLM_ENGINE_READY_TIMEOUT_S。
小结
InprocClient 和 MPClient 这两条路看起来差别巨大,但对外都长得像 EngineCoreClient——add_request 一调、get_output 一拉,前端代码完全不用关心 EngineCore 是同进程对象还是子进程 + ZMQ。同进程简单到极致,就是个对象引用;子进程那套把 ZMQ ROUTER/PULL、握手、monitor、future 表、socket ownership 全部封装好了。异步变体 AsyncMPClient 在 /client/async-mp 单独拆;EngineCore 本身见 /engine/engine-core,EngineCoreProc 的 busy loop 和 shutdown 信号见 /engine/engine-core-proc。