AsyncMPClient : greffer asyncio sur ZMQ
Responsabilités
AsyncMPClient est la variante asynchrone de MPClient, utilisée par AsyncLLM et les entrées serving. Comme SyncMPClient, elle pousse EngineCore dans un sous-processus et communique via ZMQ ROUTER pour les requêtes et PULL pour les outputs ; la différence est que le socket d'output n'est plus consommé par un thread daemon qui poll, mais directement accroché à la boucle d'événements asyncio via une task persistante qui await output_socket.recv_multipart(copy=False). Le frontal await sa requête, le sous-processus calcule et répond, la task décode et pousse dans une asyncio.Queue : toute la chaîne reste sur une seule boucle d'événements, sans commutation de thread ni passage queue.Queue inter-thread.
AsyncMPClient est implémentée avec parcimonie : elle ne surcharge que _send_input (qui renvoie un Awaitable plutôt que bloquer), un ensemble de méthodes async (call_utility_async, add_request_async, etc.) et l'essentielle _ensure_output_queue_task. La majeure partie de l'initialisation dans __init__ réutilise le handshake ZMQ et le engine monitor de MPClient.__init__ ; seul asyncio_mode=True fait basculer le contexte vers zmq.asyncio.Context et les sockets vers zmq.asyncio.Socket. Les extensions pour le parallélisme de données (DP) et EP élastique sont portées par les sous-classes DPAsyncMPClient et DPLBAsyncMPClient, abordées brièvement en fin de page.
Motivation de conception
- Boucle async mono-thread, on ne lance pas de thread daemon : le thread
process_outputs_socketdeSyncMPClientn'a pas de raison d'être en mode asynchrone —zmq.asyncio.Socket.recv_multipartest une coroutine que l'on peut directementawait, à placer dans unasyncio.create_task. La task remplace le thread, et la queue passe dequeue.Queueàasyncio.Queue. - La boucle peut ne pas être démarrée à la construction : dans
AsyncMPClient.__init__, siasyncio.get_running_loop()échoue, onpasset on diffère le lancement de la task d'output queue au premier appel de_ensure_output_queue_task(lazy task start:974-982). Cela permet de construire le client hors d'une boucle d'événements (tests, phase d'init) puis de s'y attacher paresseusement quand l'exécution commence. Ceci a un coût : entre la fin de la construction et le premieradd_request_async, si le sous-processus meurt et émetENGINE_CORE_DEAD, l'événement est perdu tant que la task n'est pas lancée — le trade-off est explicitement signalé dans le commentaire. - Éviter que la task ne tienne une référence directe sur le client :
_ensure_output_queue_taskextraitdecoder,utility_results,outputs_queueetoutput_socketcomme variables locales de closure, et utiliseweakref.ref(self)(_self_ref) pour le client lui-même. Ainsi, si le frontal lâche le client, la task ne maintient pas le client en vie via la closure (no direct ref:989-998). Si_self = _self_ref()renvoieNone, la taskreturnet se termine. - Futures utility en
asyncio.Future:_call_utility_asyncutiliseasyncio.get_running_loop().create_future()et nonconcurrent.futures.Future, de sorte queawait futuresoit directement piloté par la boucle. Le remplissage de la future passe toujours par_process_utility_outputmutualisé, qui gère les deux types de futures (_call_utility_async:1104-1116). - Notifications EEP via un call_id spécial : dans
process_outputs_socket, siutility_output.call_id == EEP_NOTIFICATION_CALL_ID, on ne touche pas à la table de futures ; on lanceasyncio.create_tasksureep_process_engine_core_notification(EEP notification path:1011-1027). L'ajout/retrait de rank en EP élastique est un « événement » et non une requête-réponse, il ne rentre pas dans le patron des futures utility. output_handlerest un hook optionnel :getattr(self.__class__, "process_engine_outputs", None)(output_handler:994-996). Les sous-classes DP (par exempleDPLBAsyncMPClient) surchargent cette méthode pour tenir des statistiques de load balancing ; la classe de baseAsyncMPClientne la définit pas, ce qui équivaut à ne pas intercepter — maisoutputs.outputs or outputs.scheduler_statsentre quand même dans la queue (入队条件:1042-1043).
Fichiers clés
AsyncMPClient.__init__:950-982— propageasyncio_mode=TrueàMPClient.__init__, créeasyncio.Queue, tente de lancer paresseusement la task si une boucle existe déjà_ensure_output_queue_task:984-1051— tire decoder/queue/utility_results/socket en locals de closure,asyncio.create_task(process_outputs_socket()), idempotentprocess_outputs_socket:1005-1047— boucle principaleawait output_socket.recv_multipart(copy=False), sépare utility / notification EEP / output ordinaireCancelledError handling:1042-1047— quand la task est annulée, pousseEngineDeadErrordans la queue pour queget_output_asynccôté appelant en soit informéget_output_async:1053-1062—await outputs_queue.get(), si on récupère une Exception on l'enveloppe dansEngineDeadErrorvia_format_exception_send_input / _send_input_message:1064-1099— renvoie unAwaitable; sans buffer c'est directement la coroutinesend_multipart, avec buffer on utilisetrack=True+add_done_callbackpour récupérer leMessageTrackeren asynchronecall_utility_async / _call_utility_async:1101-1116— génère uncall_id, enregistre uneasyncio.Future,await self._send_input_message+await futureadd_request_async / abort_requests_async:1121-1128— ajoute unclient_indexà la requête, puis appelle_ensure_output_queue_taskpour s'assurer que la task est lancéecollective_rpc_async:1188-1197— forward verscall_utility_async("collective_rpc", ...), signature alignée avec la version syncDPAsyncMPClient.__init__:1200-1240— ajoutecurrent_wave,lb_engines,first_req_send_socket,eep_scaling_cache, et lance paresseusement une task de stats update si la boucle est déjà démarrée_ensure_stats_update_task:1242-1249— lance la task d'abonnement aux stats DP, absente de la classe de baseAsyncMPClientDPLBAsyncMPClient:1380— sous-classe de load balancing interne, surchargecall_utility_asyncpour dispatcher les requêtes vers le rank le moins chargé selonlb_engines
Flux de données
Toute la chaîne tourne dans une seule boucle asyncio. Le frontal await client.add_request_async(req) → _send_input envoie les frames via send_multipart sur le ROUTER ZMQ (sans buffer, retourne directement l'awaitable de send_multipart ; avec buffer, enveloppe via track=True + done callback). Côté sous-processus, EngineCoreProc.run_busy_loop calcule puis renvoie EngineCoreOutputs sur le socket PULL. La task persistante process_outputs_socket récupère les frames via await output_socket.recv_multipart(copy=False), les décode, et les répartit par type :
async def process_outputs_socket():
try:
while True:
frames = await output_socket.recv_multipart(copy=False)
resources.validate_alive(frames)
outputs: EngineCoreOutputs = decoder.decode(frames)
if outputs.utility_output:
if (
outputs.utility_output.call_id == EEP_NOTIFICATION_CALL_ID
and notification_callback_handler is not None
):
assert _self_ref is not None
_self = _self_ref()
if not _self:
return
if outputs.utility_output.result is None:
continue
notification_data = outputs.utility_output.result.result
assert isinstance(notification_data, Sequence)
assert len(notification_data) == 2
asyncio.create_task(
notification_callback_handler(_self, notification_data)
)
else:
_process_utility_output(
outputs.utility_output, utility_results
)
continue
if output_handler is not None:
assert _self_ref is not None
_self = _self_ref()
if not _self:
# Client has been garbage collected, abort.
return
await output_handler(_self, outputs)
if outputs.outputs or outputs.scheduler_stats:
outputs_queue.put_nowait(outputs)
except Exception as e:
outputs_queue.put_nowait(e)
except asyncio.CancelledError:
outputs_queue.put_nowait(EngineDeadError())(process_outputs_socket:1005-1047). Trois branches : utility → table de futures, notification EEP → nouvelle task dédiée, outputs ordinaires → hook output_handler (s'il existe) puis file d'attente. L'output_handler est conçu pour y insérer les statistiques DP ; la classe de base ne le définit pas, les frames vont directement dans outputs_queue et le frontal les récupère via await get_output_async().
Le chemin des futures utility mérite un regard séparé, car il unifie le type de future entre sync et async :
async def _call_utility_async(
self, method: str, *args, engine: EngineIdentity
) -> Any:
call_id = uuid.uuid1().int >> 64
future = asyncio.get_running_loop().create_future()
self.utility_results[call_id] = future
message = (
EngineCoreRequestType.UTILITY.value,
*self.encoder.encode((self.client_index, call_id, method, args)),
)
await self._send_input_message(message, engine, args)
self._ensure_output_queue_task()
return await future(_call_utility_async:1104-1116). call_id utilise uuid.uuid1().int >> 64 pour prendre les 64 bits de poids fort — suffisamment grand et dispersé ; la future est enregistrée avant l'émission pour qu'une réponse plus rapide que l'enregistrement ne puisse pas survenir ; _ensure_output_queue_task() est rappelée pour s'assurer que la task tourne — sans quoi await future ne serait jamais complétée. Le remplissage se fait dans process_outputs_socket via la branche _process_utility_output(outputs.utility_output, utility_results).
Limites et échecs
- Boucle non démarrée à la construction = perte d'un signal de mort précoce :
__init__faitasyncio.get_running_loop()et en cas d'échec lance paresseusement la task (lazy:974-982), donc si le sous-processus meurt avant le premier_ensure_output_queue_task, la frameENGINE_CORE_DEADn'est jamais consommée. Le commentaire explicite ce trade-off. - En cas de cancellation de la task,
EngineDeadErrorest poussé dans la queue :except asyncio.CancelledErrorpousseEngineDeadError(CancelledError:1046-1047), de sorte queget_output_asynccôté appelant lève directement au prochainawait, sans rester suspendu silencieusement. - Auto-terminaison de la task quand le client est GC : si
_self = _self_ref()renvoieNone, onreturn(GC suicide:1037-1039), pour éviter que la task ne continue à tenir des ressources libérées. - Avec tensor buffer, renvoie une future et non une coroutine : dans
_send_input_message, silen(msg) > 3,send_multipart(track=True)renvoieasyncio.Future[zmq.MessageTracker], puisadd_done_callback(add_pending)raccroche le tracker àpending_messagesen asynchrone (track path:1091-1099). L'appelant reçoit un Awaitable qui ne se complète réellement qu'après l'envoi effectif. call_iddes notifications EEP isolé :EEP_NOTIFICATION_CALL_IDne passe pas par la tableutility_results(EEP branch:1011-1027), mais spawn une nouvelle task poureep_process_engine_core_notification, car les notifications sont des événements unidirectionnels._ensure_output_queue_taskest idempotent :if resources.output_queue_task is not None: return(idempotent:986-987).add_request_asyncetcall_utility_asyncl'appellent à chaque fois pour s'assurer que la task est lancée ; si elle l'est déjà, c'est un no-op.- Points d'extension des sous-classes DP explicites :
DPAsyncMPClientajoutecurrent_wave/lb_engines/eep_scaling_cacheet_ensure_stats_update_task;DPLBAsyncMPClientsurchargecall_utility_asyncpour faire un LB interne (DP init:1200-1240). La classe de base ne porte pas cette logique DP. - Le hook
output_handlerest cherché à la manière classmethod :getattr(self.__class__, "process_engine_outputs", None)(output_handler lookup:994-996). La sous-classe qui définit cette méthode est appelée ; la classe de base saute le hook. Il faut_self_refrécupère le client avant de l'appeler, impossible de capturerselfdirectement en closure.
Résumé
Le cœur d'AsyncMPClient est de remplacer le thread daemon + queue.Queue de SyncMPClient par une task asyncio + asyncio.Queue, le contexte ZMQ basculant vers zmq.asyncio.Context. Tout le reste — handshake ready, engine monitor, table de futures utility, MessageTracker pour empêcher le GC prématuré des tensors, propagation de ENGINE_CORE_DEAD — est hérité tel quel de la classe de base MPClient. Les extensions DP et EP élastique sont dans DPAsyncMPClient / DPLBAsyncMPClient et se branchent via le hook output_handler et la surcharge de call_utility_async. La variante synchrone est sur /client/inproc-mp ; EngineCore et EngineCoreProc eux-mêmes sur /engine/engine-core et /engine/engine-core-proc.
Voir la documentation officielle : Documentation vLLM · README