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