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 に詰め込むまでの全経路が 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 の拡張は DPAsyncMPClient、DPLBAsyncMPClient の 2 サブクラスが担い、本ページでは最後に軽く触れるだけです。
設計動機
- 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 がまだ起きていないためこのイベントは失われます。コメントはこのトレードオフを明示しています。 - 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はconcurrent.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_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 通知 / 通常出力の 3 種フレームを区分け。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で非同期にMessageTrackerを取る。call_utility_async / _call_utility_async:1101-1116—call_id生成、asyncio.Futureを登録、await self._send_input_message+await future。add_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 にリクエストを振り分け。
データフロー
経路全体が 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_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)。3 方向への分流に注意してください: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フレームが消費されません。コメントはこのトレードオフを明示的に警告しています。 - task がキャンセルされても 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 があるときは coroutine ではなく future を返す:
_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)、eep_process_engine_core_notificationを呼ぶ新 task を spawn します。通知は一方向イベントだからです。 _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を 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 を参照してください。