OpenAI Serving und AsyncLLM: asynchrone Brücke zwischen HTTP-API und Engine
Verantwortung
Nach dem Start von vllm serve läuft ein FastAPI-Server, der die OpenAI-kompatiblen Endpunkte /v1/chat/completions, /v1/completions, /v1/responses, /v1/embeddings, /pooling/*, /speech_to_text/* etc. bedient. Der Einstieg der HTTP-Schicht ist build_app(api_server.py:157), das die register_*_api_routers registriert und CORS / Authentifizierung / Ausnahme-Handler installiert. Auf der Engine-Seite ist der Einstieg AsyncLLM(async_llm.py:70), das den EngineCore (separater Prozess, Kommunikation über ZMQ) als EngineClient verpackt, den man awaiten kann, und Requests in einem asyncio-Event-Loop ausführt.
build_async_engine_client_from_engine_args(api_server.py:109) ist die Assembly-Stelle: Sie erzeugt über engine_args.create_engine_config() die VllmConfig und startet über AsyncLLM.from_vllm_config(...)(async_llm.py:202) einen AsyncLLM; intern startet EngineCoreClient.make_async_mp_client einen AsyncMPClient (Prozessübergreifend + asyncio). init_app_state(api_server.py:297) hängt engine_client, OpenAIServingModels, OnlineRenderer, ServingTokenization an app.state und initiiert danach über init_generate_state auch die Handler OpenAIServingChat, OpenAIServingCompletion, OpenAIServingResponses.
Wenn ein HTTP-Request ankommt (z. B. POST /v1/chat/completions), ruft der Router OpenAIServingChat.create_chat_completion(serving.py:233) auf. Dieser nutzt den OnlineRenderer, um die Nachrichtenliste in einen Prompt zu rendern, und erhält dann über engine_client.generate(prompt, sampling_params, request_id, ...)(serving.py:357) einen AsyncGenerator[RequestOutput], dessen gestreamte Antworten direkt zurück an HTTP geschrieben werden.
Entwurfsmotivation
- Multi-Client-Support:
AsyncLLM.__init__nimmtclient_addresses / client_count / client_index(async_llm.py:84-86), sodass ein OpenAI-Server mehrere EngineCores über DP-Ranks hinweg anbinden und Last verteilen kann. - InputProcessor / OutputProcessor getrennt:
AsyncLLMhält einenInputProcessor(der den Prompt zu einemEngineCoreRequestverarbeitet) und einenOutputProcessor(derEngineCoreOutputszurück inRequestOutputübersetzt)(async_llm.py:135-143), sodass die HTTP-Schicht nurEngineInputundRequestOutputsieht und denEngineCorenicht direkt berührt. - Hintergrund-Task output_handler:
_run_output_handler(async_llm.py:637) startet einen asyncio-Task, der vom EngineCore Outputs zieht, überoutput_processor.process_outputsin Scheiben verarbeitet undRequestOutputin denRequestOutputCollectorder jeweiligen Anfrage pusht. - Stückweises Verarbeiten gegen Event-Loop-Blockade:
output_handlerzerlegt mitVLLM_V1_OUTPUT_PROC_CHUNK_SIZE(async_llm.py:654); zwischen den Stücken wirdawait asyncio.sleep(0)ausgeführt, damit andere Anfragen yielden können. - n>1 Fan-out:
add_requestsplittet beiparams.n > 1die Anfrage in n Sub-Requests mit eigener request_id, die sich einenParentRequestund denselbenRequestOutputCollectorteilen(async_llm.py:385-397). - Automatischer Abort bei Disconnect: Wenn der HTTP-Client abbricht, wirft der
AsyncGeneratorvongenerateeinasyncio.CancelledError/GeneratorExit; der Handler fängt es ab und ruftawait self.abort(request_id, internal=True)(async_llm.py:591-596) auf, um die Anfrage aus dem EngineCore zu entfernen und GPU nicht zu verschwenden. - Sauberer Start-Fehler:
__init__verschiebt_run_output_handlerbis zur erstenadd_request(async_llm.py:370-373), sodassAsyncLLMaußerhalb des Event-Loops konstruiert werden kann und bei Startfehler direkt mit einer Exception beendet. - FastAPI teilt Engine über mehrere Endpunkte:
init_app_statepacktengine_clientinstate(api_server.py:336), sodass chat / completion / responses / pooling / speech_to_text alle denselben AsyncLLM teilen.
Schlüsseldateien
build_app:157— Konstruiert die FastAPI-App; hängt alle Router, CORS, Authentifizierung, Ausnahme-Handler ein.build_async_engine_client_from_engine_args:109— Assembly-Stelle: EngineArgs → VllmConfig → AsyncLLM.init_app_state:297— Hängt engine_client, models, renderer, serving_tokenization an app.state.AsyncLLM-Klasse:70— Asynchrone Implementierung vonEngineClient; hält InputProcessor / OutputProcessor / engine_core.AsyncLLM.from_vllm_config:202— Klassenmethoden-Factory; wählt überExecutor.get_classden Executor und konstruiert AsyncLLM.AsyncLLM.add_request:280— Einstieg: Prompt an InputProcessor übergeben, RequestOutputCollector starten,_add_requestan EngineCore schicken.AsyncLLM.generate:524—add_request+while not finished: q.get()als AsyncGenerator-Verpackung._run_output_handler:637— Hintergrund-Task:engine_core.get_output_async→output_processor.process_outputs→ Queue._add_request:400— Zwei Schritte:output_processor.add_request+engine_core.add_request_async.AsyncLLM.abort:709— Bricht eine Anfrage ab;internal=Trueist der durch Verbindungsabbruch ausgelöste interne Abort.OpenAIServingChat:106— Chat-Completions-Handler; erbt vonGenerateBaseServing.create_chat_completion:233— HTTP-Einstieg; rendert Prompt und ruft engine_client.generate.OpenAIServingCompletion:55— Handler für die Completions-API.BaseServing:29— Basisklasse aller Handler; stellt_check_model/create_error_responsebereit.
Datenfluss
Vollständiger Pfad eines POST /v1/chat/completions: FastAPI-Router → OpenAIServingChat.create_chat_completion → OnlineRenderer.render_chat() rendert die Nachrichtenliste in Prompt-Token-IDs → engine_client.generate(prompt, sampling_params, request_id, ...) liefert einen AsyncGenerator. Intern läuft generate über add_request:
# vllm/v1/engine/async_llm.py L357-L373 (vereinfacht)
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 übergibt den Prompt an InputProcessor.process_inputs, der daraus einen EngineCoreRequest macht, und ruft dann _add_request auf:
# 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)Nachdem die Anfrage im EngineCore ist, zieht der Hintergrund-Task output_handler kontinuierlich über await engine_core.get_output_async() Outputs, splittt sie in process_outputs und pusht den RequestOutput in den jeweils passenden 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)Der Aufrufer von generate läuft in einer while not finished: q.get_nowait() or await q.get()-Schleife, holt RequestOutputs und yieldet sie an die HTTP-Streaming-Antwort.
Grenzen und Fehler
- Automatischer Abort bei Client-Disconnect:
generatefängtasyncio.CancelledError/GeneratorExitab und ruftawait self.abort(request_id, internal=True)(async_llm.py:591-596) auf, um die Anfrage aus der EngineCore-Queue zu entfernen. - Engine stirbt → EngineDeadError:
generatefängtEngineDeadErrorund bricht nicht ab, sondern raise(async_llm.py:598-602), weil es nichts mehr abzubrechen gibt; die HTTP-Schicht wandelt das in 500. - n>1 Fan-out teilt Collector:
add_requestsplittet beiparams.n > 1die Anfrage in n Sub-Requests, die sich denselbenRequestOutputCollectorteilen(async_llm.py:385-397) — nach oben bleibt ein einzelner Stream. - kv_sharing_fast_prefill unterstützt kein prompt_logprobs:
add_requestwirft beikv_sharing_fast_prefill + prompt_logprobsdirektValueError(async_llm.py:305-314), weil unter diesem KV-Sharing-Pfad die Prompt-Logprobs nicht korrekt berechnet werden. - Streaming-Eingabe unterstützt kein reasoning_ended: Wenn der Input als
AsyncGeneratorübergeben wird, liefertreasoning_ended/reasoning_parser_kwargsdirektNotImplementedError(async_llm.py:316-318). - Schutz gegen Zirkelreferenz im output_handler:
_run_output_handlerspeichertengine_core/output_processor/rendererexplizit als lokale Variablen(async_llm.py:643-653), damit der Task nichtselfhält und so eine Circular Ref entsteht, die der GC nicht aufräumen kann. - request_id-Mismatch nur Warning: Wenn
add_requesteinenEngineCoreRequestübergibt, dessenrequest_idnicht übereinstimmt, wird nur gewarnt, nicht geworfen(async_llm.py:342-347); maßgeblich istEngineCoreRequest.request_id.
Zusammenfassung
OpenAI Serving ist vLLMs HTTP-Einstieg und besteht aus zwei Schichten: der FastAPI-App und den Handlern OpenAIServingChat / OpenAIServingCompletion etc. für das OpenAI-kompatible Protokoll; AsyncLLM verpackt den EngineCore als EngineClient, wirft im asyncio-Event-Loop Anfragen in den EngineCore und zieht Outputs zurück zum HTTP-Stream. Scheduling und Forward-Loop innerhalb des EngineCore siehe /engine/engine-core; wie ein HTTP-Request beim Worker berechnet wird, siehe /worker/gpu-model-runner; der letzte Schritt im Sampling siehe /sampling/sampler.