EngineCoreProc:ZMQ バックグラウンドプロセスのラッパー
役割
EngineCoreProc は EngineCore のサブクラス(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_nowait で input_queue に入れます。メインループはスケジュールの合間にキューから取り出して _handle_client_request でディスパッチ(core.py:1372-1405)します。outputs は逆方向に流れます。step_fn() が産出した EngineCoreOutputs は output_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_socketsとprocess_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_dead→raiseに進み、プロセス退出時にfinallyでengine_core.shutdown()を呼び出します(core.py:1226-1242)。
主要ファイル
EngineCoreProc class:896-910—class EngineCoreProc(EngineCore):定義 + クラス定数ENGINE_CORE_DEAD = b"ENGINE_CORE_DEAD"。EngineCoreProc.__init__:903-1009— input_queue/output_queue の起動、tensor IPC receiver、_perform_handshakes、super().__init__、2 つの IO スレッドの起動。queues + executor_fail_callback:915-919—input_queue = queue.Queue[...]、output_queue = queue.Queue[...]、executor 失敗コールバックは input_queue にEXECUTOR_FAILEDを投入。IO threads startup:974-1001—process_input_socketsとprocess_output_socketsの 2 つの daemon スレッド起動、ready_eventは input スレッドが ready になった後にのみ set。run_engine_core:1154-1224— プロセス入口、EngineCoreProcまたはDPEngineCoreProcを選択、シグナルハンドラー登録、run_busy_loop呼び出し。signal handlers:1204-1222—wakeup_engineが input_queue にWAKEUPを投入、SIGTERM/SIGINT がshutdown_stateをREQUESTEDに変更。run_busy_loop:1259-1267—while 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-1317—step_fn()呼び出し、output_queue.put_nowait(output)、post_step。forward を走らせないが scheduler に未完了があるときtime.sleep(0.001)で GIL を譲る。_handle_client_request:1372-1405—EngineCoreRequestTypeディスパッチ:WAKEUP/ADD/ABORT/UTILITY/EXECUTOR_FAILED。_send_engine_dead:1470-1482—output_queueにENGINE_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 にディスパッチします:
# 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_queue と input_queue の両方に同時投入される点に注意してください。abort はスケジューラ内で冪等(core.py:1276-1279)であるため、メインループが step の合間に aborts_queue から drain でき、abort 済みリクエストのリークを防げます。
メインループがリクエストを受け取った後、_handle_client_request が EngineCoreRequestType でディスパッチ(core.py:1372-1405)します:
# 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_coreはexcept Exceptionでengine_core._send_engine_dead()を呼び(core.py:1229-1235)、output_queueにENGINE_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_queueにEXECUTOR_FAILEDを投入し、メインループの_handle_client_requestが取り出してraise RuntimeError("Executor failed.")(core.py:1400-1401)します。プロセス全体がrun_engine_coreの except に捕まり_send_engine_deadに進みます。 - Shutdown モード:
_handle_shutdownはshutdown_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_shutdownはfailure_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 はpendingdeque にぶら下げ、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_shutdownでshutdown_stateを見て行われます。 - finally でシグナル復元:
run_engine_coreはfinallyでsignal.signal(signal.SIGTERM, signal.SIG_DFL)(core.py:1236-1242)し、シグナルハンドラーをデフォルトに戻した後engine_core.shutdown()で model_executor / scheduler / 分散状態を解放します。shutdownはgc.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 を参照してください。