OpenAI Serving と AsyncLLM:HTTP API とエンジンの非同期ブリッジ
役割
vllm serve が起動すると FastAPI サービスが走り、/v1/chat/completions、/v1/completions、/v1/responses、/v1/embeddings、/pooling/*、/speech_to_text/* などの OpenAI 互換エンドポイントをリッスンします。HTTP 層の入口は build_app(api_server.py:157)で、各 register_*_api_routers を登録し、CORS / 認証 / 例外ハンドラを取り付けます。エンジン側の入口は AsyncLLM(async_llm.py:70)で、EngineCore(独立プロセス、ZMQ で通信)を await 可能な EngineClient に包み、asyncio イベントループでリクエストを走らせます。
build_async_engine_client_from_engine_args(api_server.py:109)がアセンブリポイントです:まず engine_args.create_engine_config() で VllmConfig を得て、AsyncLLM.from_vllm_config(...)(async_llm.py:202)で AsyncLLM を起動します。内部で EngineCoreClient.make_async_mp_client が AsyncMPClient(プロセス間 + asyncio)を起動します。init_app_state(api_server.py:297)は engine_client、OpenAIServingModels、OnlineRenderer、ServingTokenization を app.state にマウントし、init_generate_state がさらに OpenAIServingChat、OpenAIServingCompletion、OpenAIServingResponses などのハンドラを構築します。
HTTP リクエストが来ると、例えば POST /v1/chat/completions はルーティングされて OpenAIServingChat.create_chat_completion(serving.py:233)を呼びます。ここで OnlineRenderer がメッセージリストを prompt にレンダリングし、engine_client.generate(prompt, sampling_params, request_id, ...)(serving.py:357)で AsyncGenerator[RequestOutput] を取得し、yield する度にストリーミングレスポンスを HTTP に書き戻します。
設計動機
- マルチクライアントサポート:
AsyncLLM.__init__はclient_addresses / client_count / client_indexを受け(async_llm.py:84-86)、1 つの OpenAI server が複数 DP rank の EngineCore に接続しロードバランスできます。 - InputProcessor / OutputProcessor の分離:
AsyncLLMはInputProcessor(prompt をEngineCoreRequestに加工)とOutputProcessor(EngineCoreOutputsをRequestOutputに翻訳)を持ち(async_llm.py:135-143)、HTTP 層にはEngineInputとRequestOutputだけが見え、EngineCoreに直接触れさせません。 - バックグラウンド output_handler タスク:
_run_output_handler(async_llm.py:637)は asyncio task を起動し、EngineCore から outputs を引き出し、output_processor.process_outputsでスライス処理し、RequestOutputを対応リクエストのRequestOutputCollectorにプッシュします。 - チャンク処理でイベントループのブロックを防止:
output_handlerはVLLM_V1_OUTPUT_PROC_CHUNK_SIZEでスライスし(async_llm.py:654)、チャンク毎にawait asyncio.sleep(0)を挟み、他のリクエストが yield できるようにします。 - n>1 fan-out:
add_requestはparams.n > 1のときリクエストを n 個のサブリクエストに分割し、それぞれ独立した request_id を持ち、ParentRequestと同じRequestOutputCollectorを共有します(async_llm.py:385-397)。 - 切断時の自動 abort:HTTP クライアントが切断すると
generateのAsyncGeneratorがasyncio.CancelledError/GeneratorExitを投げ、ハンドラが捕まえてawait self.abort(request_id, internal=True)(async_llm.py:591-596)でリクエストを EngineCore から外し、GPU を無駄使いしません。 - 起動失敗を優雅に:
__init__は_run_output_handlerを最初のadd_requestまで遅延し(async_llm.py:370-373)、AsyncLLMをイベントループ外で構築でき、起動失敗時に直接例外で抜けられます。 - FastAPI マルチエンドポイントが engine を共有:
init_app_stateはengine_clientをstateに入れ(api_server.py:336)、chat / completion / responses / pooling / speech_to_text の各ハンドラが同じ AsyncLLM を共有します。
主要ファイル
build_app:157— FastAPI app を構築し、全 router、CORS、認証、例外ハンドラを取り付けます。build_async_engine_client_from_engine_args:109— アセンブリポイント:EngineArgs → VllmConfig → AsyncLLM。init_app_state:297— engine_client、models、renderer、serving_tokenization を app.state にマウントします。AsyncLLM クラス:70—EngineClientの非同期実装、InputProcessor / OutputProcessor / engine_core を持ちます。AsyncLLM.from_vllm_config:202— クラスメソッドファクトリ、Executor.get_classでエグゼキュータを選んでから AsyncLLM を構築します。AsyncLLM.add_request:280— 入口:prompt を InputProcessor に渡し、RequestOutputCollector を起動し、_add_requestで EngineCore に投じます。AsyncLLM.generate:524—add_request+while not finished: q.get()の AsyncGenerator ラッパ。_run_output_handler:637— バックグラウンド task:engine_core.get_output_async→output_processor.process_outputs→ queue にプッシュ。_add_request:400— 2 ステップ:output_processor.add_request+engine_core.add_request_async。AsyncLLM.abort:709— リクエストをキャンセル、internal=Trueは切断で発火した内部 abort です。OpenAIServingChat:106— chat completions ハンドラ、GenerateBaseServingを継承します。create_chat_completion:233— HTTP 入口、prompt をレンダリング + engine_client.generate を呼びます。OpenAIServingCompletion:55— completions API ハンドラ。BaseServing:29— 全ハンドラの基底クラス、_check_model/create_error_responseを提供します。
データフロー
ある POST /v1/chat/completions の完全な経路:FastAPI ルーティング → OpenAIServingChat.create_chat_completion → OnlineRenderer.render_chat() がメッセージリストを prompt token ids にレンダリング → engine_client.generate(prompt, sampling_params, request_id, ...) で AsyncGenerator を取得します。generate は内部で add_request に進みます:
# vllm/v1/engine/async_llm.py L357-L373 (簡略)
generator = self.engine_client.generate(
engine_input,
sampling_params,
sub_request_id,
lora_request=lora_request,
trace_headers=trace_headers,
priority=request.priority,
data_parallel_rank=data_parallel_rank,
reasoning_ended=reasoning_ended,
reasoning_parser_kwargs={...} if parser is not None and parser.reasoning_parser is not None else None,
)add_request は prompt を InputProcessor.process_inputs に渡して EngineCoreRequest に加工してから _add_request を呼びます:
# vllm/v1/engine/async_llm.py L400-L415
async def _add_request(
self,
request: EngineCoreRequest,
prompt: str | None,
parent_req: ParentRequest | None,
index: int,
queue: RequestOutputCollector,
):
# Add the request to OutputProcessor (this process).
self.output_processor.add_request(request, prompt, parent_req, index, queue)
# Add the EngineCoreRequest to EngineCore (separate process).
await self.engine_core.add_request_async(request)
if self.log_requests:
logger.info("Added request %s.", request.request_id)リクエストが EngineCore に入った後、バックグラウンドの output_handler task はずっと await engine_core.get_output_async() で出力を引き出し、スライスして process_outputs にかけ、RequestOutput を対応 collector にプッシュします:
# vllm/v1/engine/async_llm.py L656-L683
async def output_handler():
try:
while True:
# 1) Pull EngineCoreOutputs from the EngineCore.
outputs = await engine_core.get_output_async()
num_outputs = len(outputs.outputs)
iteration_stats = (
IterationStats() if (log_stats and num_outputs) else None
)
# Split outputs into chunks of at most
# VLLM_V1_OUTPUT_PROC_CHUNK_SIZE, so that we don't block the
# event loop for too long.
engine_core_outputs = outputs.outputs
for start in range(0, num_outputs, chunk_size):
end = start + chunk_size
outputs_slice = engine_core_outputs[start:end]
# 2) Process EngineCoreOutputs.
processed_outputs = output_processor.process_outputs(
outputs_slice, outputs.timestamp, iteration_stats
)
# NOTE: RequestOutputs are pushed to their queues.
assert not processed_outputs.request_outputs
# Allow other asyncio tasks to run between chunks
if end < num_outputs:
await asyncio.sleep(0)generate の呼び出し側は while not finished: q.get_nowait() or await q.get() でループして RequestOutput を取り出し、取り出す毎に yield して HTTP ストリーミングレスポンスに送ります。
境界と失敗
- クライアント切断で自動 abort:
generateはasyncio.CancelledError/GeneratorExitを捕まえ、await self.abort(request_id, internal=True)(async_llm.py:591-596)でリクエストを EngineCore のキューから外します。 - Engine が死ぬと EngineDeadError:
generateはEngineDeadErrorを捕まえ abort せず直接 raise し(async_llm.py:598-602)、abort 対象がもうないためです。HTTP 層はこれを 500 に変換します。 - n>1 fan-out は collector を共有:
add_requestはparams.n > 1のときリクエストを n 個のサブリクエストに分割し、同じRequestOutputCollectorを共有し(async_llm.py:385-397)、上層には 1 本のストリームとして見えます。 - kv_sharing_fast_prefill は prompt_logprobs 非サポート:
add_requestはkv_sharing_fast_prefill + prompt_logprobsを検出すると直接ValueErrorを投げ(async_llm.py:305-314)、この KV 共有パスでは prompt logprobs が正確に計算できないためです。 - ストリーミング入力は reasoning_ended 非サポート:
AsyncGenerator形式のストリーミング入力を渡すときreasoning_ended/reasoning_parser_kwargsを付けると直接NotImplementedError(async_llm.py:316-318)。 - output_handler の循環参照防止:
_run_output_handlerはengine_core/output_processor/rendererを明示的にローカル変数に保存し(async_llm.py:643-653)、task がselfを持有して circular ref となり GC されないのを避けます。 - request_id 不一致は warning:
add_requestがEngineCoreRequestを渡してもrequest_idが一致しないときは warning だけで例外を投げず(async_llm.py:342-347)、EngineCoreRequest.request_idを優先します。
まとめ
OpenAI serving は vLLM の HTTP 入口で、2 層から成ります:FastAPI app + OpenAIServingChat / OpenAIServingCompletion などのハンドラが OpenAI 互換プロトコルを処理し、AsyncLLM が EngineCore を EngineClient に包んで asyncio イベントループ内でリクエストを EngineCore に投じ、出力を HTTP ストリームに引き戻します。EngineCore 内部のスケジュールと forward ループは /engine/engine-core、HTTP リクエストが最終的に worker でどう計算されるかは /worker/gpu-model-runner、サンプリングの最終ステップは /sampling/sampler を参照してください。