Skip to content

OpenAI Serving y AsyncLLM: puente asíncrono entre la API HTTP y el motor

源码版本v0.25.1

Responsabilidades

vllm serve arranca un servicio FastAPI que escucha en endpoints compatibles con OpenAI: /v1/chat/completions, /v1/completions, /v1/responses, /v1/embeddings, /pooling/*, /speech_to_text/*. La entrada de la capa HTTP es build_app(api_server.py:157), que registra cada register_*_api_routers e instala los middlewares de CORS / autenticación / manejo de excepciones. Del lado del motor, la entrada es AsyncLLM(async_llm.py:70), que envuelve a EngineCore(proceso independiente, comunicación vía ZMQ)y lo presenta como un EngineClient que se puede await, ejecutando peticiones en el event loop de asyncio.

build_async_engine_client_from_engine_args(api_server.py:109) es el punto de ensamblaje: primero engine_args.create_engine_config() para obtener VllmConfig, y luego AsyncLLM.from_vllm_config(...)(async_llm.py:202) levanta un AsyncLLM, cuyo interior EngineCoreClient.make_async_mp_client abre un AsyncMPClient(entre procesos + asyncio).init_app_state(api_server.py:297) cuelga en app.state el engine_client, OpenAIServingModels, OnlineRenderer, ServingTokenization, y luego init_generate_state termina de construir los handlers OpenAIServingChat, OpenAIServingCompletion, OpenAIServingResponses.

Cuando entra una petición HTTP, por ejemplo POST /v1/chat/completions, la ruta invoca OpenAIServingChat.create_chat_completion(serving.py:233). Este usa OnlineRenderer para renderizar la lista de mensajes en un prompt, y luego engine_client.generate(prompt, sampling_params, request_id, ...)(serving.py:357) obtiene un AsyncGenerator[RequestOutput] que va emitiendo la respuesta en streaming hacia HTTP.

Motivación de diseño

  • Soporte multi-cliente: AsyncLLM.__init__ recibe client_addresses / client_count / client_index(async_llm.py:84-86), de modo que un mismo servidor OpenAI puede conectarse a varios EngineCore por rank DP y balancear carga.
  • InputProcessor / OutputProcessor separados: AsyncLLM mantiene un InputProcessor(que convierte el prompt en EngineCoreRequest)y un OutputProcessor(que traduce EngineCoreOutputs a RequestOutput)(async_llm.py:135-143), para que la capa HTTP solo vea EngineInput y RequestOutput y no toque EngineCore directamente.
  • Tarea background output_handler: _run_output_handler(async_llm.py:637) lanza una task asyncio que saca outputs del EngineCore, los procesa por lotes en output_processor.process_outputs y empuja los RequestOutput al RequestOutputCollector de cada petición.
  • Procesamiento por chunks para no bloquear el event loop: output_handler corta con VLLM_V1_OUTPUT_PROC_CHUNK_SIZE(async_llm.py:654);entre chunk y chunk hace await asyncio.sleep(0) para que otras peticiones puedan emitir.
  • Fan-out n>1: add_request, cuando params.n > 1, descompone la petición en n subpeticiones con request_id propio, compartiendo ParentRequest y un mismo RequestOutputCollector(async_llm.py:385-397).
  • Abort automático al desconectar: si el cliente HTTP corta, el AsyncGenerator de generate recibe asyncio.CancelledError / GeneratorExit; el handler lo atrapa y ejecuta await self.abort(request_id, internal=True)(async_llm.py:591-596),retirando la petición del EngineCore para no desperdiciar GPU.
  • Fallo al arranque controlado: __init__ posterga _run_output_handler hasta la primera llamada a add_request(async_llm.py:370-373),para que AsyncLLM pueda construirse fuera del event loop y, si el arranque falla, propagar el error directamente.
  • FastAPI comparte un único engine entre endpoints: init_app_state inyecta engine_client en state(api_server.py:336);los handlers de chat / completion / responses / pooling / speech_to_text comparten el mismo AsyncLLM.

Archivos clave

Flujo de datos

Recorrido completo de un POST /v1/chat/completions: ruta FastAPI → OpenAIServingChat.create_chat_completionOnlineRenderer.render_chat() convierte la lista de mensajes en prompt token ids → engine_client.generate(prompt, sampling_params, request_id, ...) devuelve un AsyncGenerator. Por dentro, generate pasa por add_request:

python
# 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 pasa el prompt por InputProcessor.process_inputs para obtener un EngineCoreRequest, y luego llama _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)

Una vez la petición entra al EngineCore, la tarea background output_handler queda en await engine_core.get_output_async() sacando outputs, procesándolos por chunks en process_outputs y empujando los RequestOutput al collector correspondiente:

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)

El llamador de generate itera con while not finished: q.get_nowait() or await q.get() recogiendo RequestOutput, y va emitiéndolos al flujo SSE de la respuesta HTTP.

Límites y fallos

  • Abort automático al desconectar el cliente: generate atrapa asyncio.CancelledError / GeneratorExit y llama await self.abort(request_id, internal=True)(async_llm.py:591-596),retirando la petición de la cola del EngineCore.
  • Si el Engine muere, lanza EngineDeadError: generate atrapa EngineDeadError sin abortar y lo relanza directamente(async_llm.py:598-602),porque ya no hay nada que abortar; la capa HTTP lo convierte en 500.
  • Fan-out n>1 comparte collector: add_request, cuando params.n > 1, descompone la petición en n subpeticiones que comparten el mismo RequestOutputCollector(async_llm.py:385-397);el llamador sigue viendo un único flujo.
  • kv_sharing_fast_prefill no soporta prompt_logprobs: add_request, al detectar kv_sharing_fast_prefill + prompt_logprobs, lanza directamente ValueError(async_llm.py:305-314),porque en esta ruta de KV compartido los prompt logprobs no son precisos.
  • Entrada streaming no soporta reasoning_ended: si se pasa entrada como AsyncGenerator y se acompaña de reasoning_ended / reasoning_parser_kwargs, lanza NotImplementedError(async_llm.py:316-318).
  • Protección contra referencias circulares en output_handler: _run_output_handler guarda explícitamente engine_core / output_processor / renderer en variables locales(async_llm.py:643-653),para evitar que la task retenga self y crear una referencia circular que el GC no recolectaría.
  • request_id discordante solo genera warning: si add_request recibe un EngineCoreRequest cuyo request_id no cuadra, emite warning sin lanzar(async_llm.py:342-347),prevaleciendo EngineCoreRequest.request_id.

Resumen

OpenAI serving es la entrada HTTP de vLLM y se compone de dos capas: una app FastAPI con handlers OpenAIServingChat / OpenAIServingCompletion que procesan el protocolo compatible con OpenAI; y AsyncLLM, que envuelve EngineCore como EngineClient, empuja peticiones al EngineCore en el event loop de asyncio y vuelca las salidas al flujo HTTP. El bucle de planificación y forward dentro de EngineCore, en /engine/engine-core;cómo las peticiones HTTP terminan calculándose en el worker, en /worker/gpu-model-runner;el último paso de sampling, en /sampling/sampler.

Véase la documentación oficial: documentación de vLLM · README.