Skip to content

EngineCoreClient:进程内与多进程两条客户端通路

源码版本v0.25.1

职责

vLLM 的前端(LLM / AsyncLLM / serving 入口)不直接摸 EngineCore,而是经一层 client 把请求推进去、把输出拉出来。这一层就是 EngineCoreClient,它把「同步还是异步」「同进程还是子进程」这两条正交维度,统一到同一组抽象方法上。vllm/v1/engine/core_client.py 这个文件里塞下了所有实现。

最关键的两个实现是 InprocClientMPClient。前者把 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_modeasyncio_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 协议:SyncMPClientqueue.Queue 把 output socket 的帧搬到主线程,AsyncMPClientasyncio.Queue 搬,但 MPClient 共享了 encoder/decoder、ready 握手、engine monitor、utility future 注册表。差异被压在两个子类里。
  • InprocClient 的 sleep(mode="wait") 直接禁掉(InprocClient.sleep L324-L328),因为本进程模式下没有可以「等到空闲再 sleep」的并发调度器,只能 abort。这种边界在子进程模式里是支持的——子进程有自己的 busy loop 和调度器状态机。
  • utility 调用走 future 注册表:像 collective_rpcsave_sharded_state 这类需要返回值的 RPC,不能像 add_request 那样 fire-and-forget。call_utility 生成一个 call_id,把 future 塞进 self.utility_results,等 UtilityOutput 从 output socket 回来时 _process_utility_outputcall_id 取出 future 并 set_result。这套机制在 sync 和 async 两种 client 上完全一致,只是 future 类型不同(concurrent.futures.Future vs asyncio.Future)。
  • engine dead 传播:子进程意外死亡时,monitor 线程把 BackgroundResources.engine_dead 置 True,后续任何 _send_input 都会被 ensure_alive 挡住抛 EngineDeadError。output socket 那边如果收到 ENGINE_CORE_DEAD 单帧,validate_alive 也会把同一标志位拉起来,保证发请求侧和收输出侧都看得到。

关键文件

数据流

SyncMPClient 拉输出靠一个 daemon 线程在 ZMQ output_socketpoll,解码后塞进 queue.Queue。前端线程调 get_output() 就是 queue.get()——把 ZMQ 的非阻塞 socket 适配成 Python 同步阻塞接口。这段 process_outputs_socket 是整个同步通路的核心:

python
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_aliveENGINE_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_socketpoll 等子进程的 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 不支持 wait pause(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

小结

InprocClientMPClient 这两条路看起来差别巨大,但对外都长得像 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

对照官方资料:vLLM 文档 · README