Skip to content

OpenAI Serving et AsyncLLM : pont asynchrone entre l'API HTTP et le moteur

源码版本v0.25.1

Responsabilités

vllm serve démarre un service FastAPI qui expose les endpoints compatibles OpenAI : /v1/chat/completions, /v1/completions, /v1/responses, /v1/embeddings, /pooling/*, /speech_to_text/*, etc. L'entrée côté HTTP est build_app(api_server.py:157) qui enregistre les différents register_*_api_routers et installe CORS / authentification / gestionnaires d'exceptions. Côté moteur, l'entrée est AsyncLLM(async_llm.py:70) : il enveloppe le EngineCore (processus dédié, communication via ZMQ) en un EngineClient awaitable, qui s'exécute dans la boucle d'événements asyncio.

build_async_engine_client_from_engine_args(api_server.py:109) est le point d'assemblage : il appelle d'abord engine_args.create_engine_config() pour obtenir la VllmConfig, puis AsyncLLM.from_vllm_config(...)(async_llm.py:202) instancie un AsyncLLM, dont l'interne lance un AsyncMPClient via EngineCoreClient.make_async_mp_client (inter-processus + asyncio). init_app_state(api_server.py:297) attache engine_client, OpenAIServingModels, OnlineRenderer, ServingTokenization à app.state, puis init_generate_state construit également les handlers OpenAIServingChat, OpenAIServingCompletion, OpenAIServingResponses.

À l'arrivée d'une requête HTTP, par exemple POST /v1/chat/completions, le routeur appelle OpenAIServingChat.create_chat_completion(serving.py:233). Celui-ci utilise OnlineRenderer pour transformer la liste de messages en prompt, puis engine_client.generate(prompt, sampling_params, request_id, ...)(serving.py:357) récupère un AsyncGenerator[RequestOutput] qui yield au fur et à mesure les morceaux de réponse streaming vers HTTP.

Motivation de conception

  • Support multi-clients : AsyncLLM.__init__ accepte client_addresses / client_count / client_index(async_llm.py:84-86) pour qu'un serveur OpenAI puisse se connecter aux EngineCore de plusieurs ranks DP et répartir la charge.
  • Séparation InputProcessor / OutputProcessor : AsyncLLM détient un InputProcessor (transforme le prompt en EngineCoreRequest) et un OutputProcessor (traduit les EngineCoreOutputs en RequestOutput)(async_llm.py:135-143) ; la couche HTTP ne voit que EngineInput et RequestOutput, sans toucher directement au EngineCore.
  • Task output_handler en arrière-plan : _run_output_handler(async_llm.py:637) lance une tâche asyncio qui tire les sorties du EngineCore, les découpe via output_processor.process_outputs, et pousse les RequestOutput dans le RequestOutputCollector de la requête correspondante.
  • Découpage pour ne pas bloquer la boucle d'événements : output_handler découpe par VLLM_V1_OUTPUT_PROC_CHUNK_SIZE(async_llm.py:654) ; entre chaque morceau, await asyncio.sleep(0) permet aux autres requêtes de yield.
  • Fan-out n>1 : add_request éclate la requête en n sous-requêtes quand params.n > 1, chacune avec son propre request_id, mais partage le même ParentRequest et le même RequestOutputCollector(async_llm.py:385-397).
  • Abort automatique sur déconnexion : si le client HTTP se déconnecte, l'AsyncGenerator de generate reçoit une asyncio.CancelledError / GeneratorExit ; le handler la capture et appelle await self.abort(request_id, internal=True)(async_llm.py:591-596) pour retirer la requête du EngineCore et ne pas gaspiller de GPU.
  • Échec propre au démarrage : __init__ retarde _run_output_handler jusqu'au premier add_request(async_llm.py:370-373) ; AsyncLLM peut ainsi être construit hors boucle d'événements, et un échec de démarrage remonte immédiatement comme erreur.
  • Plusieurs endpoints FastAPI partagent le moteur : init_app_state place engine_client dans state(api_server.py:336) ; les handlers chat / completion / responses / pooling / speech_to_text partagent tous le même AsyncLLM.

Fichiers clés

Flux de données

Le chemin complet d'un POST /v1/chat/completions : route FastAPI → OpenAIServingChat.create_chat_completionOnlineRenderer.render_chat() transforme la liste de messages en prompt token ids → engine_client.generate(prompt, sampling_params, request_id, ...) récupère un AsyncGenerator. En interne, generate appelle add_request :

python
# vllm/v1/engine/async_llm.py L357-L373 (simplifié)
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 envoie le prompt au InputProcessor.process_inputs pour le transformer en EngineCoreRequest, puis appelle _add_request :

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)

Une fois la requête entrée dans EngineCore, la tâche d'arrière-plan output_handler ne cesse de faire await engine_core.get_output_async() pour tirer les sorties, de les découper via process_outputs, et de pousser les RequestOutput dans le collector correspondant :

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)

L'appelant de generate boucle dans while not finished: q.get_nowait() or await q.get() pour récupérer les RequestOutput, qu'il yield à la réponse HTTP streaming.

Limites et échecs

  • Abort automatique sur déconnexion client : generate capture asyncio.CancelledError / GeneratorExit et appelle await self.abort(request_id, internal=True)(async_llm.py:591-596) pour retirer la requête de la queue du EngineCore.
  • Engine mort → EngineDeadError : generate capture EngineDeadError et relance sans abort(async_llm.py:598-602) — il n'y a plus rien à aborter ; la couche HTTP le convertit en 500.
  • Fan-out n>1 partage le collector : add_request éclate en n sous-requêtes quand params.n > 1, qui partagent le même RequestOutputCollector(async_llm.py:385-397) ; l'appelant voit toujours un flux unique.
  • kv_sharing_fast_prefill incompatible avec prompt_logprobs : add_request lève ValueError si kv_sharing_fast_prefill + prompt_logprobs est détecté(async_llm.py:305-314) — les prompt logprobs ne sont pas fiables sur ce chemin de partage KV.
  • Entrée streaming incompatible avec reasoning_ended : si on passe une entrée streaming sous forme d'AsyncGenerator, transmettre reasoning_ended / reasoning_parser_kwargs lève NotImplementedError(async_llm.py:316-318).
  • Protection contre les références circulaires dans output_handler : _run_output_handler stocke explicitement engine_core / output_processor / renderer comme variables locales(async_llm.py:643-653) pour éviter que la tâche ne capture self et crée une référence circulaire que le GC ne pourrait pas collecter.
  • request_id mismatch → warning : si add_request reçoit un EngineCoreRequest dont le request_id ne correspond pas, il émet un warning sans lever(async_llm.py:342-347) ; c'est le EngineCoreRequest.request_id qui fait foi.

Résumé

OpenAI serving est l'entrée HTTP de vLLM, composée de deux couches : l'app FastAPI et les handlers OpenAIServingChat / OpenAIServingCompletion qui gèrent le protocole compatible OpenAI ; AsyncLLM enveloppe le EngineCore en EngineClient et, dans la boucle d'événements asyncio, pousse les requêtes vers le EngineCore et ramène les sorties vers le flux HTTP. La boucle de planification et de forward à l'intérieur du EngineCore est détaillée dans /engine/engine-core ; la manière dont les requêtes HTTP sont effectivement calculées par le worker dans /worker/gpu-model-runner ; la dernière étape d'échantillonnage dans /sampling/sampler.

Voir la documentation officielle : Documentation vLLM · README