EngineCoreClient:プロセス内と多プロセスの 2 つのクライアント経路
役割
vLLM のフロントエンド(LLM / AsyncLLM / serving 入口)は EngineCore を直接触らず、1 つの client レイヤーを経由してリクエストを押し込み、出力を引き出します。このレイヤーが EngineCoreClient です。「同期か非同期か」「同プロセスか子プロセスか」という 2 つの直交する次元を、同じ抽象メソッド群に統合します。vllm/v1/engine/core_client.py というファイルにすべての実装が詰め込まれています。
最も重要な 2 つの実装が InprocClient と MPClient です。前者は EngineCore を自プロセス内に直接構築し、1 ステップ呼ぶごとに 1 ステップ計算し、busy loop も ZMQ もありません。主に V0 流の LLMEngine.add_request()/step() に薄殻を被せます。後者は EngineCore を EngineCoreProc 子プロセスとして包み、ZMQ ROUTER/PULL の 2 つの socket でリクエストと出力を送受信し、フロントエンドの呼び出しスレッドは socket とだけやり取りします。MPClient は基底クラスで、具体的な実装として SyncMPClient(同期、output socket からフレームを受信するバックグラウンドスレッドを起動)と AsyncMPClient(asyncio、同じ socket を回すバックグラウンドタスクを起動)があり、後者は別ページで解説します。
EngineCoreClient.make_client は静的ファクトリで、multiprocess_mode と asyncio_mode の 2 つの bool で実装を選びます。組み合わせ(asyncio=True, multiprocess=False)は直接 NotImplementedError を投げます。EngineCore 自身は asyncio フレンドリではなく、非同期にするには多プロセスとセットで使う必要があります。
設計動機
- フロントエンドと EngineCore のライフサイクルの疎結合:同プロセスで EngineCore を呼ぶのは単純ですが、フロントエンドで複数の EngineCore を並行に走らせたり(データ並列)、弾性 EP(ランクのオンライン増減)を行うには、EngineCore を独立プロセスに押し込み、フロントエンドは socket だけを持つ必要があります。
MPClient.__init__のweakref.finalize(self, self.resources)の 1 行が核心です。構築途中で例外が出ても、バックグラウンド子プロセスと ZMQ context はBackgroundResources.__call__で兜底回収され、リークしません。 - 同期インターフェースと非同期インターフェースで同じ ZMQ プロトコルを共有:
SyncMPClientはqueue.Queueで output socket のフレームをメインスレッドに運び、AsyncMPClientはasyncio.Queueで運びますが、MPClientは encoder/decoder、ready ハンドシェイク、engine monitor、utility future レジストリを共有します。差分は 2 つのサブクラスに圧縮されています。 - InprocClient の
sleep(mode="wait")は直接拒否(InprocClient.sleepL324-L328)。これは同プロセスモードでは「アイドルになるまで待って sleep」できる並行スケジューラがいないためで、abortしか選べません。この境界は子プロセスモードではサポートされます。子プロセスには自身の busy loop とスケジューラ状態機械があるためです。 - utility 呼び出しは future レジストリ経由:
collective_rpc、save_sharded_stateのような戻り値が必要な RPC は、add_requestのような fire-and-forget では扱えません。call_utilityはcall_idを生成し、future をself.utility_resultsに詰めます。UtilityOutputが output socket から戻ってきたら_process_utility_outputがcall_idで future を取り出しset_resultします。この仕組みは sync と async の 2 種類の client で完全に一致し、future 型だけが異なります(concurrent.futures.Futurevsasyncio.Future)。 - engine dead の伝播:子プロセスが予期せず死亡したとき、monitor スレッドが
BackgroundResources.engine_deadを True にし、以後のあらゆる_send_inputはensure_aliveに弾かれてEngineDeadErrorを投げます。output socket 側でENGINE_CORE_DEAD単フレームを受け取った場合、validate_aliveも同じフラグを立て、リクエスト送信側と出力受信側の両方がそれを見られるようにします。
主要ファイル
EngineCoreClient + make_client:71-105— 抽象基底クラスと静的ファクトリ、multiprocess_mode/asyncio_modeに基づきInprocClient/SyncMPClient/AsyncMPClientを選択。make_async_mp_client:107-132— 非同期ファクトリ、data_parallel_sizeと外部 LB の有無に基づきAsyncMPClient/DPAsyncMPClient/DPLBAsyncMPClientを選択。InprocClient.__init__ / get_output:276-292— 直接EngineCoreを構築、get_outputは同期的にstep_fn()+post_step()を呼び出し。InprocClient.add_request / abort / shutdown:297-306— 直接self.engine_coreに転送、シリアライズも socket もなし。InprocClient.sleep 拒否 wait:324-328— 同プロセスモードはwaitpause 非対応、abortのみ。BackgroundResources:370-458— dataclass +__call__がweakref.finalizeのコールバック、socket を統一的に閉じ / task を停止 / engine manager を shutdown。validate_alive:454-457— 単フレームENGINE_CORE_DEAD受信でengine_deadをマークしEngineDeadErrorを送出。MPClient.__init__:480-649— ZMQ context/socket 構築、必要に応じてlaunch_core_engines、ready ハンドシェイク、monitor 起動、失敗時にはロールバック。MPClient.shutdown:651-662—detach()finalizer で engine manager 停止後、BackgroundResourcesでクリーンアップ。start_engine_core_monitor:685-712— daemon スレッドmonitor_engine_liveness、エンジン死亡時にengine_deadを立ててself.shutdown()を発火。_apply_ready_response:714-754—EngineCoreReadyResponseをデコード、max_model_len/num_gpu_blocks/block_sizeを子プロセスからフロントエンドのvllm_configに書き戻し。SyncMPClient.__init__:779-847—queue.Queue構築、process_outputs_socketdaemon スレッド起動、output socket + shutdown socket をポーリング。SyncMPClient.get_output / _send_input / call_utility:849-881— ブロッキングoutputs_queue.get()、リクエスト送信時に必要に応じてtrackMessageTracker で tensor の早期回収を防止。
データフロー
SyncMPClient は daemon スレッドで ZMQ output_socket を poll し、デコード後に queue.Queue に詰め込むことで出力を引っ張ります。フロントエンドスレッドが get_output() を呼ぶのは queue.get() と同義で、ZMQ の非ブロッキング socket を Python 同期ブロッキングインターフェースに適合させます。この process_outputs_socket が同期経路全体の核心です:
def process_outputs_socket():
assert isinstance(out_socket, zmq.Socket)
shutdown_socket = ctx.socket(zmq.PAIR)
try:
shutdown_socket.bind(shutdown_path)
poller = zmq.Poller()
poller.register(shutdown_socket, zmq.POLLIN)
poller.register(out_socket, zmq.POLLIN)
while True:
socks = poller.poll()
if not socks:
continue
if len(socks) == 2 or socks[0][0] == shutdown_socket:
# shutdown signal, exit thread.
break
frames = out_socket.recv_multipart(copy=False)
resources.validate_alive(frames)
outputs: EngineCoreOutputs = decoder.decode(frames)
if outputs.utility_output:
_process_utility_output(outputs.utility_output, utility_results)
else:
outputs_queue.put_nowait(outputs)
except Exception as e:
outputs_queue.put_nowait(e)
finally:
# Close sockets.
shutdown_socket.close(linger=0)
out_socket.close(linger=0)(SyncMPClient.process_outputs_socket:808-836)。3 つの詳細に注意してください。1 つ目、shutdown は専用の zmq.PAIR inproc socket を shutdown_path に bind します。BackgroundResources.__call__ が閉じる際にこの socket に空バイトを送ってスレッドを退出させ、スレッドが out_socket を待ち続けないようにします。2 つ目、validate_alive は ENGINE_CORE_DEAD 単フレームを EngineDeadError に変換します。3 つ目、utility 返信(collective_rpc の戻り値など)と通常の outputs は同じフレームで流れ、前者は utility_results[call_id] の future に直接流し込まれ、後者はキューに入ります。
リクエスト送信方向はさらに単純で、SyncMPClient._send_input は直接 (identity, request_type, *encoder.encode(request)) の複数フレームを ROUTER に送ります。tensor buffer があるときは track=True で MessageTracker を取得し、pending_messages にぶら下げて Python 側が backing tensor を早期 GC するのを防ぎます(_send_input:861-873)。
ハンドシェイク段階は MPClient.__init__ で(等待 ready:615-634):各 engine_rank について 2 バイトの little-endian identity を生成し、input_socket で poll して子プロセスの EngineCoreReadyResponse を待ち、タイムアウトは VLLM_ENGINE_READY_TIMEOUT_S で抑えられます。_apply_ready_response が子プロセスで算出された num_gpu_blocks、整列済みの block_size、自動フィットされた max_model_len をフロントエンドの vllm_config に書き戻します。これがフロントエンドがスケジュールと流量制限に使える本当の容量です。
境界と失敗
asyncio_mode=True, multiprocess_mode=Falseは直接拒否(make_client:91-95)。EngineCore 内部は async-friendly ではなく、この組み合わせは実装されていません。- 構築失敗時は子プロセスを回収:
MPClient.__init__はtry/finally+successフラグで、失敗時にself._finalizer()を呼び出してBackgroundResources.__call__を発火し、既にlaunch_core_enginesで起動した子プロセスをまとめて停止(__init__ finally:499-649)します。 - InprocClient は
waitpause をサポートしない(InprocClient.sleep:324-328)。同プロセスモードではアイドルを待つ主体がおらず、現在のリクエストを abort してから sleep することしかできません。 - engine core 子プロセスの死亡は両方向から感知される:送信側は
ensure_alive(ensure_alive:670-672)で、受信側はvalidate_aliveがENGINE_CORE_DEAD単フレームを受け取る(validate_alive:454-457)ことで。両方ともBackgroundResources.engine_deadを立てます。 - utility future の InvalidStateError 許容:
_process_utility_outputはasyncio.InvalidStateErrorを捕捉(_process_utility_output:769-776)し、呼び出し側 task がキャンセルされて future が既に cancelled になっているシナリオをカバーします。 - output socket のクローズは context.term の前:
BackgroundResources.__call__は明示的にclose_socketsしてから後続(sync case:437-452)を行わないと、ZMQ context の終了が詰まります。 - SyncMPClient の output_socket はスレッドが閉じる:構築の最後に
self.resources.output_socket = None(socket ownership 转交:846-847)とし、BackgroundResourcesと daemon スレッドが同時に閉じるのを防ぎます。 - ready ハンドシェイクタイムアウトは読みやすいエラー:
Timed out waiting for engine core processes to start. This is often caused by slow weight loading for large models.(ready timeout:622-630)——ユーザーにVLLM_ENGINE_READY_TIMEOUT_Sの調整を直接伝えます。
まとめ
InprocClient と MPClient の 2 経路は見かけ上の差が大きいですが、外部からはどちらも EngineCoreClient に見えます。add_request を呼び、get_output を引くだけで、フロントエンドコードは EngineCore が同プロセスのオブジェクトか子プロセス + ZMQ かを意識しません。同プロセスは極限までシンプルで単なるオブジェクト参照です。子プロセス側は ZMQ ROUTER/PULL、ハンドシェイク、monitor、future レジストリ、socket ownership をすべてカプセル化します。非同期版 AsyncMPClient は /client/async-mp で別途取り上げます。EngineCore 自身は /engine/engine-core、EngineCoreProc の busy loop と shutdown シグナルは /engine/engine-core-proc を参照してください。