OpenAI Serving y AsyncLLM: puente asíncrono entre la API HTTP y el motor
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__recibeclient_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:
AsyncLLMmantiene unInputProcessor(que convierte el prompt enEngineCoreRequest)y unOutputProcessor(que traduceEngineCoreOutputsaRequestOutput)(async_llm.py:135-143), para que la capa HTTP solo veaEngineInputyRequestOutputy no toqueEngineCoredirectamente. - Tarea background output_handler:
_run_output_handler(async_llm.py:637) lanza una task asyncio que saca outputs del EngineCore, los procesa por lotes enoutput_processor.process_outputsy empuja losRequestOutputalRequestOutputCollectorde cada petición. - Procesamiento por chunks para no bloquear el event loop:
output_handlercorta conVLLM_V1_OUTPUT_PROC_CHUNK_SIZE(async_llm.py:654);entre chunk y chunk haceawait asyncio.sleep(0)para que otras peticiones puedan emitir. - Fan-out n>1:
add_request, cuandoparams.n > 1, descompone la petición en n subpeticiones con request_id propio, compartiendoParentRequesty un mismoRequestOutputCollector(async_llm.py:385-397). - Abort automático al desconectar: si el cliente HTTP corta, el
AsyncGeneratordegeneraterecibeasyncio.CancelledError/GeneratorExit; el handler lo atrapa y ejecutaawait 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_handlerhasta la primera llamada aadd_request(async_llm.py:370-373),para queAsyncLLMpueda construirse fuera del event loop y, si el arranque falla, propagar el error directamente. - FastAPI comparte un único engine entre endpoints:
init_app_stateinyectaengine_clientenstate(api_server.py:336);los handlers de chat / completion / responses / pooling / speech_to_text comparten el mismo AsyncLLM.
Archivos clave
build_app:157— construye la app FastAPI, monta routers, CORS, autenticación y handlers de excepciones.build_async_engine_client_from_engine_args:109— punto de ensamblaje: EngineArgs → VllmConfig → AsyncLLM.init_app_state:297— cuelga engine_client, models, renderer y serving_tokenization en app.state.AsyncLLM 类:70— implementación asíncrona deEngineClient, mantiene InputProcessor / OutputProcessor / engine_core.AsyncLLM.from_vllm_config:202— método fábrica que llamaExecutor.get_classpara elegir el ejecutor y luego construye AsyncLLM.AsyncLLM.add_request:280— entrada: pasa el prompt por InputProcessor, crea un RequestOutputCollector y_add_requestlo envía al EngineCore.AsyncLLM.generate:524— envolturaadd_request+while not finished: q.get()como AsyncGenerator._run_output_handler:637— tarea background:engine_core.get_output_async→output_processor.process_outputs→ push a la cola._add_request:400— dos pasos:output_processor.add_request+engine_core.add_request_async.AsyncLLM.abort:709— cancela la petición;internal=Trueindica abort interno por desconexión.OpenAIServingChat:106— handler de chat completions, hereda deGenerateBaseServing.create_chat_completion:233— entrada HTTP, renderiza el prompt y llama a engine_client.generate.OpenAIServingCompletion:55— handler de la API de completions.BaseServing:29— clase base de todos los handlers, ofrece_check_model/create_error_response.
Flujo de datos
Recorrido completo de un POST /v1/chat/completions: ruta FastAPI → OpenAIServingChat.create_chat_completion → OnlineRenderer.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:
# 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:
# 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:
# 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:
generateatrapaasyncio.CancelledError/GeneratorExity llamaawait 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:
generateatrapaEngineDeadErrorsin 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, cuandoparams.n > 1, descompone la petición en n subpeticiones que comparten el mismoRequestOutputCollector(async_llm.py:385-397);el llamador sigue viendo un único flujo. - kv_sharing_fast_prefill no soporta prompt_logprobs:
add_request, al detectarkv_sharing_fast_prefill + prompt_logprobs, lanza directamenteValueError(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
AsyncGeneratory se acompaña dereasoning_ended/reasoning_parser_kwargs, lanzaNotImplementedError(async_llm.py:316-318). - Protección contra referencias circulares en output_handler:
_run_output_handlerguarda explícitamenteengine_core/output_processor/rendereren variables locales(async_llm.py:643-653),para evitar que la task retengaselfy crear una referencia circular que el GC no recolectaría. - request_id discordante solo genera warning: si
add_requestrecibe unEngineCoreRequestcuyorequest_idno cuadra, emite warning sin lanzar(async_llm.py:342-347),prevaleciendoEngineCoreRequest.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.