Skip to content

OpenAI Serving und AsyncLLM: asynchrone Brücke zwischen HTTP-API und Engine

源码版本v0.25.1

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__ nimmt client_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: AsyncLLM hält einen InputProcessor (der den Prompt zu einem EngineCoreRequest verarbeitet) und einen OutputProcessor (der EngineCoreOutputs zurück in RequestOutput übersetzt)(async_llm.py:135-143), sodass die HTTP-Schicht nur EngineInput und RequestOutput sieht und den EngineCore nicht direkt berührt.
  • Hintergrund-Task output_handler: _run_output_handler(async_llm.py:637) startet einen asyncio-Task, der vom EngineCore Outputs zieht, über output_processor.process_outputs in Scheiben verarbeitet und RequestOutput in den RequestOutputCollector der jeweiligen Anfrage pusht.
  • Stückweises Verarbeiten gegen Event-Loop-Blockade: output_handler zerlegt mit VLLM_V1_OUTPUT_PROC_CHUNK_SIZE(async_llm.py:654); zwischen den Stücken wird await asyncio.sleep(0) ausgeführt, damit andere Anfragen yielden können.
  • n>1 Fan-out: add_request splittet bei params.n > 1 die Anfrage in n Sub-Requests mit eigener request_id, die sich einen ParentRequest und denselben RequestOutputCollector teilen(async_llm.py:385-397).
  • Automatischer Abort bei Disconnect: Wenn der HTTP-Client abbricht, wirft der AsyncGenerator von generate ein asyncio.CancelledError / GeneratorExit; der Handler fängt es ab und ruft await 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_handler bis zur ersten add_request(async_llm.py:370-373), sodass AsyncLLM außerhalb des Event-Loops konstruiert werden kann und bei Startfehler direkt mit einer Exception beendet.
  • FastAPI teilt Engine über mehrere Endpunkte: init_app_state packt engine_client in state(api_server.py:336), sodass chat / completion / responses / pooling / speech_to_text alle denselben AsyncLLM teilen.

Schlüsseldateien

Datenfluss

Vollständiger Pfad eines POST /v1/chat/completions: FastAPI-Router → OpenAIServingChat.create_chat_completionOnlineRenderer.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:

python
# 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:

python
# 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:

python
# 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: generate fängt asyncio.CancelledError / GeneratorExit ab und ruft await self.abort(request_id, internal=True)(async_llm.py:591-596) auf, um die Anfrage aus der EngineCore-Queue zu entfernen.
  • Engine stirbt → EngineDeadError: generate fängt EngineDeadError und 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_request splittet bei params.n > 1 die Anfrage in n Sub-Requests, die sich denselben RequestOutputCollector teilen(async_llm.py:385-397) — nach oben bleibt ein einzelner Stream.
  • kv_sharing_fast_prefill unterstützt kein prompt_logprobs: add_request wirft bei kv_sharing_fast_prefill + prompt_logprobs direkt ValueError(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, liefert reasoning_ended / reasoning_parser_kwargs direkt NotImplementedError(async_llm.py:316-318).
  • Schutz gegen Zirkelreferenz im output_handler: _run_output_handler speichert engine_core / output_processor / renderer explizit als lokale Variablen(async_llm.py:643-653), damit der Task nicht self hält und so eine Circular Ref entsteht, die der GC nicht aufräumen kann.
  • request_id-Mismatch nur Warning: Wenn add_request einen EngineCoreRequest übergibt, dessen request_id nicht übereinstimmt, wird nur gewarnt, nicht geworfen(async_llm.py:342-347); maßgeblich ist EngineCoreRequest.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.

Siehe offizielle Dokumentation: vLLM 文档 · README.