Skip to content

AsyncMPClient : greffer asyncio sur ZMQ

源码版本v0.25.1

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_socket de SyncMPClient n'a pas de raison d'être en mode asynchrone — zmq.asyncio.Socket.recv_multipart est une coroutine que l'on peut directement await, à placer dans un asyncio.create_task. La task remplace le thread, et la queue passe de queue.Queue à asyncio.Queue.
  • La boucle peut ne pas être démarrée à la construction : dans AsyncMPClient.__init__, si asyncio.get_running_loop() échoue, on pass et 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 premier add_request_async, si le sous-processus meurt et émet ENGINE_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_task extrait decoder, utility_results, outputs_queue et output_socket comme variables locales de closure, et utilise weakref.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() renvoie None, la task return et se termine.
  • Futures utility en asyncio.Future : _call_utility_async utilise asyncio.get_running_loop().create_future() et non concurrent.futures.Future, de sorte que await future soit directement piloté par la boucle. Le remplissage de la future passe toujours par _process_utility_output mutualisé, 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, si utility_output.call_id == EEP_NOTIFICATION_CALL_ID, on ne touche pas à la table de futures ; on lance asyncio.create_task sur eep_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_handler est un hook optionnel : getattr(self.__class__, "process_engine_outputs", None) (output_handler:994-996). Les sous-classes DP (par exemple DPLBAsyncMPClient) surchargent cette méthode pour tenir des statistiques de load balancing ; la classe de base AsyncMPClient ne la définit pas, ce qui équivaut à ne pas intercepter — mais outputs.outputs or outputs.scheduler_stats entre quand même dans la queue (入队条件:1042-1043).

Fichiers clés

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 :

python
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 :

python
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__ fait asyncio.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 frame ENGINE_CORE_DEAD n'est jamais consommée. Le commentaire explicite ce trade-off.
  • En cas de cancellation de la task, EngineDeadError est poussé dans la queue : except asyncio.CancelledError pousse EngineDeadError (CancelledError:1046-1047), de sorte que get_output_async côté appelant lève directement au prochain await, sans rester suspendu silencieusement.
  • Auto-terminaison de la task quand le client est GC : si _self = _self_ref() renvoie None, on return (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, si len(msg) > 3, send_multipart(track=True) renvoie asyncio.Future[zmq.MessageTracker], puis add_done_callback(add_pending) raccroche le tracker à pending_messages en 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_id des notifications EEP isolé : EEP_NOTIFICATION_CALL_ID ne passe pas par la table utility_results (EEP branch:1011-1027), mais spawn une nouvelle task pour eep_process_engine_core_notification, car les notifications sont des événements unidirectionnels.
  • _ensure_output_queue_task est idempotent : if resources.output_queue_task is not None: return (idempotent:986-987). add_request_async et call_utility_async l'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 : DPAsyncMPClient ajoute current_wave / lb_engines / eep_scaling_cache et _ensure_stats_update_task ; DPLBAsyncMPClient surcharge call_utility_async pour faire un LB interne (DP init:1200-1240). La classe de base ne porte pas cette logique DP.
  • Le hook output_handler est 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_ref récupère le client avant de l'appeler, impossible de capturer self directement 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