Skip to content

EngineCoreProc:ZMQ バックグラウンドプロセスのラッパー

源码版本v0.25.1

役割

EngineCoreProcEngineCore のサブクラス(core.py:896-897)で、これを独立したバックグラウンドプロセス (background process) で丸ごと走らせます。docstring は非常に直接的で ZMQ-wrapper for running EngineCore in background process. と書かれています。行うのは 3 つだけです。プロセス起動時にフロントエンドとハンドシェイクして ZMQ socket アドレスを取得し、2 つの IO スレッドを起動して socket と 2 つの queue.Queue を橋渡し(core.py:915-916)、その後 run_busy_loop をメインスレッドで回し、input_queue からリクエストを取り出し → step_fn() でスケジュールと前向き計算を実行 → outputs を output_queue に詰め込む(core.py:1259-1267)ことを繰り返します。

このレイヤーは v1 の多プロセスアーキテクチャの鍵です。フロントエンドプロセス(SyncMPClient または AsyncMPClient)は ZMQ で EngineCoreRequest をシリアライズして送り込み、EngineCoreProc の input スレッドが受け取って put_nowaitinput_queue に入れます。メインループはスケジュールの合間にキューから取り出して _handle_client_request でディスパッチ(core.py:1372-1405)します。outputs は逆方向に流れます。step_fn() が産出した EngineCoreOutputsoutput_queue に放り込まれ、output スレッドがキューから取り出してシリアライズし ZMQ でフロントエンドに送り返します。ZMQ socket IO は送受信時に GIL を解放するため、この 2 つの IO スレッドは GPU 前向き計算と本当に重ねられ(core.py:974-1001)、シリアライズ/デシリアライズも forward の裏に隠せます。

設計動機

  • プロセス分離:EngineCoreProc は独立プロセスで走るため、フロントエンドの Python がクラッシュしても既に GPU に入った worker を巻き込まず、worker の OOM も API プロセスを道連れにしません。これこそ v1 が EngineCore を外に出した核心的な理由です。v0 の LLMEngine は同プロセス内ですべてを行い、1 箇所のクラッシュが全部を道連れにしていました。
  • ZMQ + queue.Queue の 2 段バッファ:ZMQ socket をメインループで直接使わず、中間に queue.Queue を挟みます(core.py:915-916)。socket IO は IO スレッドに置き、メインループは queue.get(block=...) だけで済ませます。メインループが socket recv でブロックされたり ZMQ 内部ロックに引きずられたりしないようにしつつ、socket IO が GIL を解放するとメインループは即座に GIL を取り戻して forward を走らせられます。
  • 2 つの IO スレッド:process_input_socketsprocess_output_sockets(core.py:980-1001)はそれぞれ独立し、入力スレッドはデシリアライズ + put_nowait を担当、出力スレッドは get + シリアライズ + send_multipart を担当します。両者は互いにブロックしません。
  • 起動ハンドシェイク (handshake):コンストラクタで _perform_handshakes(core.py:932-938)が EngineZmqAddresses のセットを取得し、その後 input スレッドが ROUTER/DEALER 上で EngineCoreReadyResponse を送信してフロントエンドに「起きたよ、max_model_len、num_gpu_blocks、kv_cache_size_tokens はこれだよ」と伝えます(core.py:1525-1546)。フロントエンドは MPClient.__init__ でこの ready メッセージを待ち、届かなければ起動失敗とみなします。
  • run_engine_core プロセス入口:この関数は multiprocessing.Process(target=run_engine_core, ...) の target(core.py:1154-1224)です。EngineCoreProc をインスタンス化、SIGTERM/SIGINT を登録、最後に run_busy_loop を呼び出します。例外はすべて except Exception_send_engine_deadraise に進み、プロセス退出時に finallyengine_core.shutdown() を呼び出します(core.py:1226-1242)。

主要ファイル

  • EngineCoreProc class:896-910class EngineCoreProc(EngineCore): 定義 + クラス定数 ENGINE_CORE_DEAD = b"ENGINE_CORE_DEAD"
  • EngineCoreProc.__init__:903-1009 — input_queue/output_queue の起動、tensor IPC receiver、_perform_handshakessuper().__init__、2 つの IO スレッドの起動。
  • queues + executor_fail_callback:915-919input_queue = queue.Queue[...]output_queue = queue.Queue[...]、executor 失敗コールバックは input_queue に EXECUTOR_FAILED を投入。
  • IO threads startup:974-1001process_input_socketsprocess_output_sockets の 2 つの daemon スレッド起動、ready_event は input スレッドが ready になった後にのみ set。
  • run_engine_core:1154-1224 — プロセス入口、EngineCoreProc または DPEngineCoreProc を選択、シグナルハンドラー登録、run_busy_loop 呼び出し。
  • signal handlers:1204-1222wakeup_engine が input_queue に WAKEUP を投入、SIGTERM/SIGINT が shutdown_stateREQUESTED に変更。
  • run_busy_loop:1259-1267while self._handle_shutdown(): _process_input_queue(); _process_engine_step()、shutdown 後に raise SystemExit
  • _process_input_queue:1269-1298 — 仕事がないときは input_queue.get(block=True) でブロック、仕事があるときは非ブロックで drain。
  • _process_engine_step:1300-1317step_fn() 呼び出し、output_queue.put_nowait(output)post_step。forward を走らせないが scheduler に未完了があるとき time.sleep(0.001) で GIL を譲る。
  • _handle_client_request:1372-1405EngineCoreRequestType ディスパッチ:WAKEUP/ADD/ABORT/UTILITY/EXECUTOR_FAILED。
  • _send_engine_dead:1470-1482output_queueENGINE_CORE_DEAD を投入し、output スレッドがメッセージを送出するのを待ってから退出。
  • process_input_sockets:1484-1587 — input スレッド本体、zmq.Poller で DEALER + オプションの coord XSUB を監視、デシリアライズ後に put_nowait で input_queue へ。
  • process_output_sockets:1589-1654 — output スレッド本体、output_queue.get() で outputs を取り出し、encoder.encode_into + send_multipart(copy=False, track=True) でゼロコピー送信。

データフロー

リクエストは ZMQ input socket から入ってきて、まず input スレッドでデシリアライズされ input_queue に詰め込まれ、メインループが _process_input_queue で取り出して _handle_client_request にディスパッチします:

python
# vllm/v1/engine/core.py L1556-L1587
while True:
    for input_socket, _ in poller.poll():
        # (RequestType, RequestData)
        type_frame, *data_frames = input_socket.recv_multipart(copy=False)
        # NOTE(yongji): ignore READY message sent by DP coordinator
        # that is used to notify newly started engines
        if type_frame.buffer == b"READY":
            assert input_socket == coord_socket
            continue
        request_type = EngineCoreRequestType(bytes(type_frame.buffer))

        # Deserialize the request data.
        request: Any
        if request_type == EngineCoreRequestType.ADD:
            req: EngineCoreRequest = add_request_decoder.decode(data_frames)
            try:
                request = self.preprocess_add_request(req)
            except Exception:
                self._handle_request_preproc_error(req)
                continue
        else:
            request = generic_decoder.decode(data_frames)

            if request_type == EngineCoreRequestType.ABORT:
                # Aborts are added to *both* queues, allows us to eagerly
                # process aborts while also ensuring ordering in the input
                # queue to avoid leaking requests. This is ok because
                # aborting in the scheduler is idempotent.
                self.aborts_queue.put_nowait(request)

        # Push to input queue for core busy loop.
        self.input_queue.put_nowait((request_type, request))

ABORT は aborts_queueinput_queue の両方に同時投入される点に注意してください。abort はスケジューラ内で冪等(core.py:1276-1279)であるため、メインループが step の合間に aborts_queue から drain でき、abort 済みリクエストのリークを防げます。

メインループがリクエストを受け取った後、_handle_client_requestEngineCoreRequestType でディスパッチ(core.py:1372-1405)します:

python
# vllm/v1/engine/core.py L1377-L1401
if request_type == EngineCoreRequestType.WAKEUP:
    return
elif request_type == EngineCoreRequestType.ADD:
    req, request_wave = request
    if self._reject_add_in_shutdown(req):
        return
    self.add_request(req, request_wave)
elif request_type == EngineCoreRequestType.ABORT:
    self.abort_requests(request)
elif request_type == EngineCoreRequestType.UTILITY:
    client_idx, call_id, method_name, args = request
    if self._reject_utility_in_shutdown(client_idx, call_id, method_name):
        return
    output = UtilityOutput(call_id)
    # Lazily look-up utility method so that failure will be handled/returned.
    get_result = lambda: (
        (method := getattr(self, method_name))
        and method(*self._convert_msgspec_args(method, args))
    )
    enqueue_output = lambda out: self.output_queue.put_nowait(
        (client_idx, EngineCoreOutputs(utility_output=out))
    )
    self._invoke_utility_method(method_name, get_result, output, enqueue_output)
elif request_type == EngineCoreRequestType.EXECUTOR_FAILED:
    raise RuntimeError("Executor failed.")

UTILITY は汎用 RPC チャンネルです。フロントエンドがメソッド名 + msgspec args を送り、EngineCoreProc が getattr でリフレクション呼び出しし、結果は output_queue 経由で返送します。_invoke_utility_method は戻り値が Future の場合も処理でき(core.py:1440-1446)、非同期ツール完了後に output を返送します。

境界と失敗

  • プロセスクラッシュ通知:run_engine_coreexcept Exceptionengine_core._send_engine_dead() を呼び(core.py:1229-1235)、output_queueENGINE_CORE_DEAD バイト列を投入し(core.py:1470-1474)、output スレッドがこの特殊値を受け取ると全 output socket にブロードキャスト(core.py:1625-1628)します。フロントエンドの MPClient は受け取った後にエラーを投げ、無限待機しません。
  • Executor 失敗コールバック:__init__executor_fail_callback クロージャを model_executor に登録(core.py:917-919)します。worker 内部でエラー時にコールバックが input_queueEXECUTOR_FAILED を投入し、メインループの _handle_client_request が取り出して raise RuntimeError("Executor failed.")(core.py:1400-1401)します。プロセス全体が run_engine_core の except に捕まり _send_engine_dead に進みます。
  • Shutdown モード:_handle_shutdownshutdown_state(core.py:1324-1360)を見て、REQUESTED 状態で shutdown_timeout==0 なら abort で即座に全リクエストを殺し、0 以外なら drain で全リクエスト完了を待ちます。shutdown 中に新しく入ってきた ADD と UTILITY は明示的に拒否(core.py:1407-1432)し、_reject_add_in_shutdown はクライアントに abort output を返し、_reject_utility_in_shutdownfailure_message="Server shutting down" を付けた UtilityOutput を返します。
  • input スレッド死亡時のエラー報告:__init__ready_event.wait(timeout=10) で input スレッドの ready をループ待機(core.py:1003-1009)し、10 秒たっても ready にならずスレッドが死んでいれば raise RuntimeError("Input socket thread died during startup") します。まだ ready でなくても生きていれば、DP Coordinator の READY メッセージを待ち続けます。
  • ゼロコピー送信 + MessageTracker:output スレッドは send_multipart(buffers, copy=False, track=True)(core.py:1646-1654)で MessageTracker を返し、送信未完了の tracker と outputs 参照と buffer は pending deque にぶら下げ、GC が未送信の backing buffer を回収しないようにします。reuse_buffers は最大 len(sockets) + 1 個の再利用バッファに制限し、メモリの無限増大を防ぎます。
  • WAKEUP は何もしない操作:シグナルハンドラーが wakeup_engine で input_queue に WAKEUP を投入し、input_queue.get(block=True) でブロックされているメインループを起こすだけ(core.py:1377-1378)です。WAKEUP 自体は何もせず、本当の shutdown 処理は _handle_shutdownshutdown_state を見て行われます。
  • finally でシグナル復元:run_engine_corefinallysignal.signal(signal.SIGTERM, signal.SIG_DFL)(core.py:1236-1242)し、シグナルハンドラーをデフォルトに戻した後 engine_core.shutdown() で model_executor / scheduler / 分散状態を解放します。shutdowngc.unfreeze()(core.py:652-656)も行い、起動時に gc.freeze() で凍結した startup heap を再び GC に渡します。さもなくば同プロセスモード(単体テスト)で GPU メモリがリークします。

まとめ

EngineCoreProc は v1 多プロセスアーキテクチャの耐力壁です。EngineCore を継承し、独立プロセス内で ZMQ + 2 つの queue.Queue + 2 つの IO スレッドを走らせ、socket IO と GPU 前向き計算を本当に重ね合わせます。run_engine_core はプロセス入口で、インスタンス化、シグナル登録、例外兜底、shutdown を担当します。すべてのフロントエンド(SyncMPClient / AsyncMPClient)とバックエンド (EngineCore) の通信はこのレイヤーを経由します。さらに中を見るには /engine/engine-core でエンジンコア自身の組み立て、/engine/llm-engine でフロントエンド薄殻、クライアントの 3 つの実装は /client/inproc-mp、スケジューラ内部は /scheduler/scheduler を参照してください。

公式資料:vLLM 文档 · README