Skip to content

AsyncMPClient: envolver ZMQ con un event loop asyncio

源码版本v0.25.1

Responsabilidades

AsyncMPClient es la variante asíncrona de MPClient, pensada para AsyncLLM y las entradas de serving. Igual que SyncMPClient, lleva EngineCore a un subproceso y se comunica mediante un ZMQ ROUTER para enviar peticiones y un PULL para recibir salidas; la diferencia está en que el socket de salida ya no se encuesta desde un hilo demonio, sino que se monta directamente sobre el event loop asyncio: una tarea persistente hace await output_socket.recv_multipart(copy=False) para leer los frames. Así, el await del frontend, el cómputo del subproceso y el reenvío a una asyncio.Queue viven todos en el mismo event loop, sin cambio de hilo ni queue.Queue para mover datos entre hilos.

La propia AsyncMPClient está implementada con restricción: solo sobrescribe _send_input (devuelve un Awaitable en vez de bloquear), un grupo de métodos async (call_utility_async/add_request_async, etc.) y el crucial _ensure_output_queue_task. La mayoría de la lógica de inicialización de __init__ reutiliza el handshake ZMQ y el engine monitor de MPClient.__init__; lo único que cambia es que asyncio_mode=True hace que el context pase a ser zmq.asyncio.Context y el socket a zmq.asyncio.Socket. Las extensiones de paralelismo de datos (DP) y de EP elástico las aportan las subclases DPAsyncMPClient y DPLBAsyncMPClient, que se mencionan brevemente al final.

Motivación de diseño

  • Event loop de un solo hilo, no se puede abrir otro demonio: el hilo demonio process_outputs_socket de SyncMPClient sobra en el modo asíncrono —zmq.asyncio.Socket.recv_multipart ya es una coroutine y se puede await directamente, basta con empaquetarla en un asyncio.create_task. La tarea sustituye al hilo, y la cola pasa de queue.Queue a asyncio.Queue.
  • El loop puede no estar activo durante la construcción: en AsyncMPClient.__init__, si asyncio.get_running_loop() falla, se hace pass y se pospone el arranque de la tarea de output queue hasta la primera llamada a _ensure_output_queue_task(lazy task start:974-982). Así el client puede construirse fuera del event loop (típicamente en tests o fases de inicialización) y acoplarse después de forma perezosa. El trade-off está documentado: entre la construcción y la primera llamada a add_request_async, si el subproceso muere y emite ENGINE_CORE_DEAD, la tarea aún no está activa y el evento se pierde.
  • La tarea no debe retener al client: _ensure_output_queue_task saca decoder, utility_results, outputs_queue y output_socket como variables locales del closure, y el propio client se referencia con weakref.ref(self) (_self_ref). Así, si el frontend suelta al client, la tarea no lo mantiene vivo por el closure(no direct ref:989-998). Si _self = _self_ref() devuelve None, la tarea hace return y se suicida.
  • utility future con asyncio.Future: _call_utility_async usa asyncio.get_running_loop().create_future() en vez de concurrent.futures.Future, para que await future pueda ser impulsado directamente por el event loop. El rellenado del future sigue pasando por la función compartida _process_utility_output, compatible con ambos tipos de future(_call_utility_async:1104-1116).
  • Las notificaciones EEP usan un call_id especial: dentro de process_outputs_socket, si utility_output.call_id == EEP_NOTIFICATION_CALL_ID, no se consulta la tabla de futures, sino que se lanza un asyncio.create_task que dispara el callback eep_process_engine_core_notification(ruta de notificación EEP:1011-1027). Las notificaciones de añadir o quitar rank en EP elástico son «eventos», no «petición-respuesta», y no encajan en el mecanismo de utility future.
  • output_handler es un hook opcional: getattr(self.__class__, "process_engine_outputs", None)(output_handler:994-996). Las subclases DP (como DPLBAsyncMPClient) lo sobrescriben para llevar estadísticas de balanceo de carga; la clase base AsyncMPClient no lo define, lo que equivale a no interceptar —pero outputs.outputs or outputs.scheduler_stats sigue entrando en la cola(condición de encolado:1042-1043).

Archivos clave

Flujo de datos

Toda la cadena corre dentro de un único loop asyncio. El frontend hace await client.add_request_async(req)_send_input envía los frames con send_multipart al ROUTER de ZMQ (sin buffer devuelve directamente el awaitable de send_multipart; con buffer lo envuelve con track=True + done callback). En el subproceso, EngineCoreProc.run_busy_loop computa y envía de vuelta los EngineCoreOutputs al socket PULL. La tarea persistente process_outputs_socket recibe los frames en await output_socket.recv_multipart(copy=False), los decodifica y los reparte según su tipo:

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). Obsérvense las tres rutas: la respuesta utility pasa por la tabla de futures, la notificación EEP lanza una tarea nueva y la salida normal pasa por el hook output_handler (si existe) antes de entrar en la cola. El diseño de output_handler está pensado para que las subclases DP cuelguen estadísticas —la clase base no lo define, los frames van directamente a outputs_queue, y el frontend los recoge con await get_output_async().

Conviene mirar por separado la ruta del future en las llamadas utility, porque ahí se unifican los tipos de future síncrono y asíncrono:

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 usa los 64 bits altos de uuid.uuid1().int, suficientemente grande y disperso; el future se registra antes de enviar la petición, para que la respuesta no llegue antes que el registro; la llamada a _ensure_output_queue_task() asegura que la tarea esté activa —de lo contrario await future no se rellenaría nunca. El rellenado lo dispara la rama _process_utility_output(outputs.utility_output, utility_results) dentro de process_outputs_socket.

Límites y fallos

  • Si el loop no está activo durante la construcción, se pierde una señal de muerte temprana: en __init__, el fallo de asyncio.get_running_loop() inicia la tarea de forma perezosa(lazy:974-982); si el subproceso muere antes de la primera llamada a _ensure_output_queue_task, el frame ENGINE_CORE_DEAD no se consume. El comentario advierte explícitamente de este trade-off.
  • La tarea cancelada también inyecta un EngineDeadError: el except asyncio.CancelledError mete EngineDeadError en la cola(CancelledError:1046-1047); así, el await get_output_async del frontend lanza una excepción en vez de colgarse en silencio.
  • La tarea se suicida si el client entra en GC: cuando _self = _self_ref() devuelve None, la tarea hace return(GC suicide:1037-1039), para no sostener recursos ya liberados.
  • Con tensor buffer se devuelve un future, no una coroutine: en _send_input_message, si len(msg) > 3, send_multipart(track=True) devuelve un asyncio.Future[zmq.MessageTracker] y, con add_done_callback(add_pending), el tracker se cuelga de forma asíncrona en pending_messages(track path:1091-1099). El llamador recibe un Awaitable, y el envío solo se completa de verdad tras el await.
  • call_id de notificación EEP es independiente: EEP_NOTIFICATION_CALL_ID no pasa por la tabla utility_results(EEP branch:1011-1027); en su lugar se lanza una tarea nueva que invoca eep_process_engine_core_notification, porque la notificación es un evento unidireccional.
  • _ensure_output_queue_task es idempotente: if resources.output_queue_task is not None: return(idempotente:986-987). add_request_async y call_utility_async lo llaman cada vez para asegurar que la tarea esté activa; si ya está, no hace nada.
  • Los puntos de extensión de las subclases DP están claramente delimitados: DPAsyncMPClient añade current_wave / lb_engines / eep_scaling_cache y _ensure_stats_update_task; DPLBAsyncMPClient sobrescribe call_utility_async para hacer LB interno(DP init:1200-1240). La clase base no carga esta lógica de DP.
  • El hook output_handler se busca al estilo classmethod: getattr(self.__class__, "process_engine_outputs", None)(output_handler lookup:994-996). Si una subclase lo define, se invoca; si no, se omite. Hay que sacar _self_ref antes de invocarlo y no capturar self en el closure directamente.

Resumen

El núcleo de AsyncMPClient es reemplazar el hilo demonio + queue.Queue de SyncMPClient por una tarea asyncio + asyncio.Queue, y conmutar el context ZMQ a zmq.asyncio.Context. Todo lo demás —handshake de ready, engine monitor, tabla de futures utility, MessageTracker para evitar el GC prematuro de tensores, propagación de ENGINE_CORE_DEAD— se hereda de MPClient. Las extensiones de DP y EP elástico viven en DPAsyncMPClient / DPLBAsyncMPClient y se enganchan mediante el hook output_handler y la sobrescritura de call_utility_async. La variante síncrona está en /client/inproc-mp; EngineCore y EngineCoreProc en sí están en /engine/engine-core y /engine/engine-core-proc.

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