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。