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 に詰め込むまでの全経路が 1 つのイベントループ内に収まり、スレッド切り替えも queue.Queue のクロススレッド運搬もありません。

AsyncMPClient 自体の実装は控えめです。_send_input(ブロッキングではなく Awaitable を返す)、call_utility_async/add_request_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 の 2 サブクラスが担い、本ページでは最後に軽く触れるだけです。

設計動機

  • async loop はシングルスレッド、これ以上 daemon スレッドを開かない:SyncMPClientprocess_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 がまだ起きていないためこのイベントは失われます。コメントはこのトレードオフを明示しています。
  • 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_asyncconcurrent.futures.Future ではなく asyncio.get_running_loop().create_future() を使います。これにより await future がイベントループで直接駆動できます。future の書き戻しは引き続き共有の _process_utility_output 経由で、この関数は両方の future 型に対応(_call_utility_async:1104-1116)します。
  • EEP 通知は特殊な call_id:process_outputs_socketutility_output.call_id == EEP_NOTIFICATION_CALL_ID の場合、future レジストリを経由せず、asyncio.create_taskeep_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)。

主要ファイル

データフロー

経路全体が 1 つの 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)。3 方向への分流に注意してください: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 フレームが消費されません。コメントはこのトレードオフを明示的に警告しています。
  • task がキャンセルされても EngineDeadError を詰め込む:except asyncio.CancelledErrorEngineDeadError をキューに流し込み(CancelledError:1046-1047)、get_output_async 側が await して直接エラーを投げるようにします。知らぬ間にハングアップしません。
  • client が GC されたら task は自殺:_self = _self_ref()None のとき return し(GC 自杀:1037-1039)、task が既に解放されたリソースを持ち続けないようにします。
  • tensor buffer があるときは coroutine ではなく future を返す:_send_input_messagelen(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_IDutility_results レジストリを経由せず(EEP branch:1011-1027)、eep_process_engine_core_notification を呼ぶ新 task を spawn します。通知は一方向イベントだからです。
  • _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 を追加し、DPLBAsyncMPClientcall_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 を closure 捕獲できない点に注意してください。

まとめ

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