AsyncMPClient:给 ZMQ 套上 asyncio 事件循环
职责
AsyncMPClient 是 MPClient 的异步变体,给 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 的扩展由 DPAsyncMPClient、DPLBAsyncMPClient 两个子类承担,本页最后点到为止。
设计动机
- async loop 单线程,不能再开 daemon 线程:
SyncMPClient那条process_outputs_socketdaemon 线程在异步模式下没必要——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_task把decoder、utility_results、outputs_queue、output_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_async用asyncio.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)。
关键文件
AsyncMPClient.__init__:950-982— 透传asyncio_mode=True给MPClient.__init__,建asyncio.Queue,尝试在已有 loop 里 lazy 起 task_ensure_output_queue_task:984-1051— 把 decoder/queue/utility_results/socket 拉成 closure 局部变量,asyncio.create_task(process_outputs_socket()),幂等process_outputs_socket:1005-1047—await output_socket.recv_multipart(copy=False)主循环,区分 utility / EEP 通知 / 普通输出三类帧CancelledError 处理:1042-1047— task 被取消时往队列塞EngineDeadError,让get_output_async那侧也能感知get_output_async:1053-1062—await outputs_queue.get(),若拿到 Exception 用_format_exception包成EngineDeadError_send_input / _send_input_message:1064-1099— 返回Awaitable,无 buffer 时直接返回send_multipartcoroutine,有 buffer 时用track=True+add_done_callback异步拿MessageTrackercall_utility_async / _call_utility_async:1101-1116— 生成call_id、注册asyncio.Future、await self._send_input_message+await futureadd_request_async / abort_requests_async:1121-1128— 给请求打上client_index,发完顺手_ensure_output_queue_task保证 task 起来collective_rpc_async:1188-1197— 透传到call_utility_async("collective_rpc", ...),和 sync 版参数对齐DPAsyncMPClient.__init__:1200-1240— 多了current_wave、lb_engines、first_req_send_socket、eep_scaling_cache,在 loop 已起时 lazy 起 stats update task_ensure_stats_update_task:1242-1249— 起 DP 统计订阅 task,基类AsyncMPClient没有DPLBAsyncMPClient:1380— 内部负载均衡子类,覆写call_utility_async把请求按lb_engines派到最闲的 rank
数据流
整条链路在一个 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_socket 在 await output_socket.recv_multipart(copy=False) 处拿到帧,解码后按类型分流:
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 类型统一起来了:
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_id 用 uuid.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.CancelledError把EngineDeadError灌进队列(CancelledError:1046-1047),这样get_output_async那侧 await 出来直接抛错,不会无声挂死。 - client 被 GC 时 task 自杀:
_self = _self_ref()取到None时return(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_async、call_utility_async每次都会调它,确保 task 起来;已起就 no-op。- DP 子类扩展点明确:
DPAsyncMPClient加current_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。