OpenAI Serving et AsyncLLM : pont asynchrone entre l'API HTTP et le moteur
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__accepteclient_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 :
AsyncLLMdétient unInputProcessor(transforme le prompt enEngineCoreRequest) et unOutputProcessor(traduit lesEngineCoreOutputsenRequestOutput)(async_llm.py:135-143) ; la couche HTTP ne voit queEngineInputetRequestOutput, sans toucher directement auEngineCore. - 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 viaoutput_processor.process_outputs, et pousse lesRequestOutputdans leRequestOutputCollectorde la requête correspondante. - Découpage pour ne pas bloquer la boucle d'événements :
output_handlerdécoupe parVLLM_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 quandparams.n > 1, chacune avec son propre request_id, mais partage le mêmeParentRequestet le mêmeRequestOutputCollector(async_llm.py:385-397). - Abort automatique sur déconnexion : si le client HTTP se déconnecte, l'
AsyncGeneratordegeneratereçoit uneasyncio.CancelledError/GeneratorExit; le handler la capture et appelleawait 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_handlerjusqu'au premieradd_request(async_llm.py:370-373) ;AsyncLLMpeut 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_stateplaceengine_clientdansstate(api_server.py:336) ; les handlers chat / completion / responses / pooling / speech_to_text partagent tous le même AsyncLLM.
Fichiers clés
build_app:157— construit l'app FastAPI, attache tous les routers, CORS, authentification, gestionnaires d'exceptions.build_async_engine_client_from_engine_args:109— point d'assemblage : EngineArgs → VllmConfig → AsyncLLM.init_app_state:297— attache engine_client, models, renderer, serving_tokenization à app.state.AsyncLLM 类:70— implémentation asynchrone deEngineClient, porte InputProcessor / OutputProcessor / engine_core.AsyncLLM.from_vllm_config:202— méthode usine, appelleExecutor.get_classpour sélectionner l'exécuteur puis construit AsyncLLM.AsyncLLM.add_request:280— entrée : envoie le prompt au InputProcessor, crée le RequestOutputCollector, délègue à_add_requestvers EngineCore.AsyncLLM.generate:524— wrapper AsyncGenerator autour deadd_request+while not finished: q.get()._run_output_handler:637— tâche d'arrière-plan :engine_core.get_output_async→output_processor.process_outputs→ push dans la queue._add_request:400— en deux étapes :output_processor.add_request+engine_core.add_request_async.AsyncLLM.abort:709— annule une requête,internal=Truecorrespond à un abort interne déclenché par une déconnexion.OpenAIServingChat:106— handler chat completions, hérite deGenerateBaseServing.create_chat_completion:233— entrée HTTP, render du prompt + appel à engine_client.generate.OpenAIServingCompletion:55— handler de l'API completions.BaseServing:29— classe de base de tous les handlers, fournit_check_model/create_error_response.
Flux de données
Le chemin complet d'un POST /v1/chat/completions : route FastAPI → OpenAIServingChat.create_chat_completion → OnlineRenderer.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 :
# 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 :
# 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 :
# 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 :
generatecaptureasyncio.CancelledError/GeneratorExitet appelleawait self.abort(request_id, internal=True)(async_llm.py:591-596) pour retirer la requête de la queue du EngineCore. - Engine mort → EngineDeadError :
generatecaptureEngineDeadErroret 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 quandparams.n > 1, qui partagent le mêmeRequestOutputCollector(async_llm.py:385-397) ; l'appelant voit toujours un flux unique. - kv_sharing_fast_prefill incompatible avec prompt_logprobs :
add_requestlèveValueErrorsikv_sharing_fast_prefill + prompt_logprobsest 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, transmettrereasoning_ended/reasoning_parser_kwargslèveNotImplementedError(async_llm.py:316-318). - Protection contre les références circulaires dans output_handler :
_run_output_handlerstocke explicitementengine_core/output_processor/renderercomme variables locales(async_llm.py:643-653) pour éviter que la tâche ne captureselfet crée une référence circulaire que le GC ne pourrait pas collecter. - request_id mismatch → warning : si
add_requestreçoit unEngineCoreRequestdont lerequest_idne correspond pas, il émet un warning sans lever(async_llm.py:342-347) ; c'est leEngineCoreRequest.request_idqui 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