Skip to content

EngineCoreClient:プロセス内と多プロセスの 2 つのクライアント経路

源码版本v0.25.1

役割

vLLM のフロントエンド(LLM / AsyncLLM / serving 入口)は EngineCore を直接触らず、1 つの client レイヤーを経由してリクエストを押し込み、出力を引き出します。このレイヤーが EngineCoreClient です。「同期か非同期か」「同プロセスか子プロセスか」という 2 つの直交する次元を、同じ抽象メソッド群に統合します。vllm/v1/engine/core_client.py というファイルにすべての実装が詰め込まれています。

最も重要な 2 つの実装が InprocClientMPClient です。前者は 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_modeasyncio_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 プロトコルを共有:SyncMPClientqueue.Queue で output socket のフレームをメインスレッドに運び、AsyncMPClientasyncio.Queue で運びますが、MPClient は encoder/decoder、ready ハンドシェイク、engine monitor、utility future レジストリを共有します。差分は 2 つのサブクラスに圧縮されています。
  • InprocClient の sleep(mode="wait") は直接拒否(InprocClient.sleep L324-L328)。これは同プロセスモードでは「アイドルになるまで待って sleep」できる並行スケジューラがいないためで、abort しか選べません。この境界は子プロセスモードではサポートされます。子プロセスには自身の busy loop とスケジューラ状態機械があるためです。
  • utility 呼び出しは future レジストリ経由:collective_rpcsave_sharded_state のような戻り値が必要な RPC は、add_request のような fire-and-forget では扱えません。call_utilitycall_id を生成し、future を self.utility_results に詰めます。UtilityOutput が output socket から戻ってきたら _process_utility_outputcall_id で future を取り出し set_result します。この仕組みは sync と async の 2 種類の client で完全に一致し、future 型だけが異なります(concurrent.futures.Future vs asyncio.Future)。
  • engine dead の伝播:子プロセスが予期せず死亡したとき、monitor スレッドが BackgroundResources.engine_dead を True にし、以後のあらゆる _send_inputensure_alive に弾かれて EngineDeadError を投げます。output socket 側で ENGINE_CORE_DEAD 単フレームを受け取った場合、validate_alive も同じフラグを立て、リクエスト送信側と出力受信側の両方がそれを見られるようにします。

主要ファイル

データフロー

SyncMPClient は daemon スレッドで ZMQ output_socketpoll し、デコード後に queue.Queue に詰め込むことで出力を引っ張ります。フロントエンドスレッドが get_output() を呼ぶのは queue.get() と同義で、ZMQ の非ブロッキング socket を Python 同期ブロッキングインターフェースに適合させます。この process_outputs_socket が同期経路全体の核心です:

python
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_aliveENGINE_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=TrueMessageTracker を取得し、pending_messages にぶら下げて Python 側が backing tensor を早期 GC するのを防ぎます(_send_input:861-873)。

ハンドシェイク段階は MPClient.__init__ で(等待 ready:615-634):各 engine_rank について 2 バイトの little-endian identity を生成し、input_socketpoll して子プロセスの 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 は wait pause をサポートしない(InprocClient.sleep:324-328)。同プロセスモードではアイドルを待つ主体がおらず、現在のリクエストを abort してから sleep することしかできません。
  • engine core 子プロセスの死亡は両方向から感知される:送信側は ensure_alive(ensure_alive:670-672)で、受信側は validate_aliveENGINE_CORE_DEAD 単フレームを受け取る(validate_alive:454-457)ことで。両方とも BackgroundResources.engine_dead を立てます。
  • utility future の InvalidStateError 許容:_process_utility_outputasyncio.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 の調整を直接伝えます。

まとめ

InprocClientMPClient の 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 を参照してください。

公式資料:vLLM 文档 · README