Skip to content

AsyncMPClient:给 ZMQ 套上 asyncio 事件循环

源码版本v0.25.1

职责

AsyncMPClientMPClient 的异步变体,给 AsyncLLM 和 serving 入口用。和 SyncMPClient 一样,它把 EngineCore 推到子进程里跑,通过 ZMQ ROUTER 发请求、PULL 收输出;差别在于 output socket 不再由一个 daemon 线程 poll,而是直接挂在 asyncio 事件循环上,用一个常驻 task await output_socket.recv_multipart(copy=False) 拉帧。这样一来,前端 await 出去、子进程算完回包、task 解码后塞进 asyncio.Queue,整个链路全程在一个事件循环里,没有线程切换、没有 queue.Queue 跨线程搬运。

AsyncMPClient 本身实现得很克制:它只覆盖了 _send_input(返回 Awaitable 而非阻塞)、call_utility_async/add_request_async 等一组 async 方法,以及关键的 _ensure_output_queue_task__init__ 里大多数初始化逻辑都复用 MPClient.__init__ 的 ZMQ 握手和 engine monitor 那一套,只是 asyncio_mode=True 让 context 换成 zmq.asyncio.Context,socket 换成 zmq.asyncio.Socket。数据并行(DP)和弹性 EP 的扩展由 DPAsyncMPClientDPLBAsyncMPClient 两个子类承担,本页最后点到为止。

设计动机

  • async loop 单线程,不能再开 daemon 线程:SyncMPClient 那条 process_outputs_socket daemon 线程在异步模式下没必要——zmq.asyncio.Socket.recv_multipart 本身就是 coroutine,可以直接 await,塞进一个 asyncio.create_task 就完事。线程被 task 替代,queue 也从 queue.Queue 换成 asyncio.Queue
  • 构造时 loop 可能还没起:AsyncMPClient.__init__asyncio.get_running_loop() 失败就 pass,把 output queue task 的启动延迟到第一次 _ensure_output_queue_task 被调用(lazy task start:974-982)。这是为了让 client 可以在事件循环外被构造(比如测试或初始化阶段),真正跑起来时再 lazy attach。但这有个代价——构造完到第一次 add_request_async 之间,如果子进程死亡发回 ENGINE_CORE_DEAD,因为 task 还没起,这个事件会丢;注释里明确点出这个 trade-off。
  • 避免 task 直接持有 client 引用:_ensure_output_queue_taskdecoderutility_resultsoutputs_queueoutput_socket 都拿出来当 closure 局部变量,client 本体用 weakref.ref(self)(_self_ref)。这样如果前端把 client 丢了,task 不会因为 closure 引用撑着 client 不被 GC(no direct ref:989-998)。_self = _self_ref() 取出来若是 None 直接 return,task 自杀。
  • utility future 用 asyncio.Future:_call_utility_asyncasyncio.get_running_loop().create_future() 而不是 concurrent.futures.Future,这样 await future 才能直接被事件循环驱动。future 的回填仍走共享的 _process_utility_output,这个函数对两种 future 都兼容(_call_utility_async:1104-1116)。
  • EEP 通知走特殊 call_id:process_outputs_socket 里如果 utility_output.call_id == EEP_NOTIFICATION_CALL_ID,不走 future 表,而是 asyncio.create_task 触发 eep_process_engine_core_notification 回调(EEP 通知路径:1011-1027)。弹性 EP 的 rank 增减通知是「事件」而不是「请求-响应」,不能套 utility future 那套。
  • output_handler 是可选钩子:getattr(self.__class__, "process_engine_outputs", None)(output_handler:994-996)。DP 子类(如 DPLBAsyncMPClient)会覆写这个方法做负载均衡统计,基类 AsyncMPClient 没有这个方法,等于不拦——但 outputs.outputs or outputs.scheduler_stats 仍然会进队列(入队条件:1042-1043)。

关键文件

数据流

整条链路在一个 asyncio loop 里跑。前端 await client.add_request_async(req)_send_input 把帧 send_multipart 到 ZMQ ROUTER(无 buffer 时直接返回 send_multipart 的 awaitable,有 buffer 时包一层 track=True + done callback)。子进程那边 EngineCoreProc.run_busy_loop 算完把 EngineCoreOutputs 发回 PULL socket。常驻 task process_outputs_socketawait output_socket.recv_multipart(copy=False) 处拿到帧,解码后按类型分流:

python
async def process_outputs_socket():
    try:
        while True:
            frames = await output_socket.recv_multipart(copy=False)
            resources.validate_alive(frames)
            outputs: EngineCoreOutputs = decoder.decode(frames)
            if outputs.utility_output:
                if (
                    outputs.utility_output.call_id == EEP_NOTIFICATION_CALL_ID
                    and notification_callback_handler is not None
                ):
                    assert _self_ref is not None
                    _self = _self_ref()
                    if not _self:
                        return
                    if outputs.utility_output.result is None:
                        continue
                    notification_data = outputs.utility_output.result.result
                    assert isinstance(notification_data, Sequence)
                    assert len(notification_data) == 2
                    asyncio.create_task(
                        notification_callback_handler(_self, notification_data)
                    )
                else:
                    _process_utility_output(
                        outputs.utility_output, utility_results
                    )
                continue

            if output_handler is not None:
                assert _self_ref is not None
                _self = _self_ref()
                if not _self:
                    # Client has been garbage collected, abort.
                    return
                await output_handler(_self, outputs)

            if outputs.outputs or outputs.scheduler_stats:
                outputs_queue.put_nowait(outputs)
    except Exception as e:
        outputs_queue.put_nowait(e)
    except asyncio.CancelledError:
        outputs_queue.put_nowait(EngineDeadError())

(process_outputs_socket:1005-1047)。注意三路分流:utility 回包走 future 表、EEP 通知起一个新 task 处理、普通输出走 output_handler 钩子(若存在)再入队列。output_handler 的设计是给 DP 子类插统计用的——基类没有这个方法,帧直接进 outputs_queue,前端 await get_output_async() 拿到的就是这些帧。

utility 调用的 future 路径要单独看一下,因为它把 sync 和 async 的 future 类型统一起来了:

python
async def _call_utility_async(
    self, method: str, *args, engine: EngineIdentity
) -> Any:
    call_id = uuid.uuid1().int >> 64
    future = asyncio.get_running_loop().create_future()
    self.utility_results[call_id] = future
    message = (
        EngineCoreRequestType.UTILITY.value,
        *self.encoder.encode((self.client_index, call_id, method, args)),
    )
    await self._send_input_message(message, engine, args)
    self._ensure_output_queue_task()
    return await future

(_call_utility_async:1104-1116)。call_iduuid.uuid1().int >> 64 取高 64 位,够大够稀疏;future 注册完才发请求,避免回包比注册快;_ensure_output_queue_task() 再调一次保证 task 起来——否则 await future 永远不会被回填。回填发生在 process_outputs_socket 里调 _process_utility_output(outputs.utility_output, utility_results) 那一支。

边界与失败

  • 构造时 loop 未起会丢早期死亡信号:__init__asyncio.get_running_loop() 失败就 lazy 起 task(lazy:974-982),意味着如果子进程在第一次 _ensure_output_queue_task 之前死亡,ENGINE_CORE_DEAD 帧不会被消费。注释明确警示这个 trade-off。
  • task 被 cancel 时也塞 EngineDeadError:except asyncio.CancelledErrorEngineDeadError 灌进队列(CancelledError:1046-1047),这样 get_output_async 那侧 await 出来直接抛错,不会无声挂死。
  • client 被 GC 时 task 自杀:_self = _self_ref() 取到 Nonereturn(GC 自杀:1037-1039),避免 task 持续持有已释放资源。
  • 有 tensor buffer 时返回 future 而非 coroutine:_send_input_message 里如果 len(msg) > 3,send_multipart(track=True) 返回 asyncio.Future[zmq.MessageTracker],再 add_done_callback(add_pending) 异步把 tracker 挂到 pending_messages(track path:1091-1099)。调用方拿到的是一个 Awaitable,await 后才真正完成发送。
  • EEP 通知 call_id 独立:EEP_NOTIFICATION_CALL_ID 不走 utility_results 表(EEP branch:1011-1027),而是 spawn 一个新 task 调 eep_process_engine_core_notification,因为通知是单向事件。
  • _ensure_output_queue_task 幂等:if resources.output_queue_task is not None: return(幂等:986-987)。add_request_asynccall_utility_async 每次都会调它,确保 task 起来;已起就 no-op。
  • DP 子类扩展点明确:DPAsyncMPClientcurrent_wave / lb_engines / eep_scaling_cache_ensure_stats_update_task;DPLBAsyncMPClient 覆写 call_utility_async 做内部 LB(DP init:1200-1240)。基类不背这些 DP 逻辑。
  • output_handler 钩子是 classmethod 风格的查找:getattr(self.__class__, "process_engine_outputs", None)(output_handler 查找:994-996)。子类定义这个方法就会被调用,基类不定义就跳过。注意它要 _self_ref 取到 client 才能调,不能直接闭包捕获 self

小结

AsyncMPClient 的核心是把 SyncMPClient 那条 daemon 线程 + queue.Queue 换成 asyncio task + asyncio.Queue,ZMQ context 也跟着切到 zmq.asyncio.Context。其余一切——ready 握手、engine monitor、utility future 表、MessageTracker 防 tensor 提前 GC、ENGINE_CORE_DEAD 传播——全部从 MPClient 基类继承下来。DP 和弹性 EP 的扩展在 DPAsyncMPClient / DPLBAsyncMPClient 里,通过 output_handler 钩子和覆写 call_utility_async 接入。同步变体见 /client/inproc-mp,EngineCore 和 EngineCoreProc 本身见 /engine/engine-core/engine/engine-core-proc

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