AsyncMPClient: envolver ZMQ con un event loop asyncio
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_socketdeSyncMPClientsobra en el modo asíncrono —zmq.asyncio.Socket.recv_multipartya es una coroutine y se puedeawaitdirectamente, basta con empaquetarla en unasyncio.create_task. La tarea sustituye al hilo, y la cola pasa dequeue.Queueaasyncio.Queue. - El loop puede no estar activo durante la construcción: en
AsyncMPClient.__init__, siasyncio.get_running_loop()falla, se hacepassy 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 aadd_request_async, si el subproceso muere y emiteENGINE_CORE_DEAD, la tarea aún no está activa y el evento se pierde. - La tarea no debe retener al client:
_ensure_output_queue_tasksacadecoder,utility_results,outputs_queueyoutput_socketcomo variables locales del closure, y el propio client se referencia conweakref.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()devuelveNone, la tarea hacereturny se suicida. - utility future con asyncio.Future:
_call_utility_asyncusaasyncio.get_running_loop().create_future()en vez deconcurrent.futures.Future, para queawait futurepueda 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, siutility_output.call_id == EEP_NOTIFICATION_CALL_ID, no se consulta la tabla de futures, sino que se lanza unasyncio.create_taskque dispara el callbackeep_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 (comoDPLBAsyncMPClient) lo sobrescriben para llevar estadísticas de balanceo de carga; la clase baseAsyncMPClientno lo define, lo que equivale a no interceptar —perooutputs.outputs or outputs.scheduler_statssigue entrando en la cola(condición de encolado:1042-1043).
Archivos clave
AsyncMPClient.__init__:950-982— pasaasyncio_mode=TrueaMPClient.__init__, crea laasyncio.Queuee intenta arrancar la tarea de forma perezosa si hay loop_ensure_output_queue_task:984-1051— saca decoder/queue/utility_results/socket como locales del closure, lanzaasyncio.create_task(process_outputs_socket()), idempotenteprocess_outputs_socket:1005-1047— bucle principal conawait output_socket.recv_multipart(copy=False), distingue tres tipos de frame: utility / notificación EEP / salida normalCancelledError handling:1042-1047— cuando la tarea se cancela, inyecta unEngineDeadErroren la cola para queget_output_asynctambién lo percibaget_output_async:1053-1062—await outputs_queue.get(); si recibe una excepción la envuelve enEngineDeadErrorcon_format_exception_send_input / _send_input_message:1064-1099— devuelve unAwaitable; sin buffer retorna la coroutinesend_multipart, con buffer usatrack=True+add_done_callbackpara obtener unMessageTrackerde forma asíncronacall_utility_async / _call_utility_async:1101-1116— genera uncall_id, registra unasyncio.Future, haceawait self._send_input_message+await futureadd_request_async / abort_requests_async:1121-1128— etiqueta la petición conclient_indexy, tras enviar, llama a_ensure_output_queue_taskpara asegurar que la tarea esté activacollective_rpc_async:1188-1197— delega encall_utility_async("collective_rpc", ...), alineando los parámetros con la versión síncronaDPAsyncMPClient.__init__:1200-1240— añadecurrent_wave,lb_engines,first_req_send_socket,eep_scaling_cache; con el loop activo arranca perezosamente la tarea de stats_ensure_stats_update_task:1242-1249— arranca la tarea de suscripción a estadísticas de DP; la clase baseAsyncMPClientno la tieneDPLBAsyncMPClient:1380— subclase interna de balanceo de carga, sobrescribecall_utility_asyncpara repartir las peticiones al rank más libre segúnlb_engines
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:
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:
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 deasyncio.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 frameENGINE_CORE_DEADno se consume. El comentario advierte explícitamente de este trade-off. - La tarea cancelada también inyecta un EngineDeadError: el
except asyncio.CancelledErrormeteEngineDeadErroren la cola(CancelledError:1046-1047); así, elawait get_output_asyncdel 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()devuelveNone, la tarea hacereturn(GC suicide:1037-1039), para no sostener recursos ya liberados. - Con tensor buffer se devuelve un future, no una coroutine: en
_send_input_message, silen(msg) > 3,send_multipart(track=True)devuelve unasyncio.Future[zmq.MessageTracker]y, conadd_done_callback(add_pending), el tracker se cuelga de forma asíncrona enpending_messages(track path:1091-1099). El llamador recibe un Awaitable, y el envío solo se completa de verdad tras elawait. - call_id de notificación EEP es independiente:
EEP_NOTIFICATION_CALL_IDno pasa por la tablautility_results(EEP branch:1011-1027); en su lugar se lanza una tarea nueva que invocaeep_process_engine_core_notification, porque la notificación es un evento unidireccional. _ensure_output_queue_taskes idempotente:if resources.output_queue_task is not None: return(idempotente:986-987).add_request_asyncycall_utility_asynclo 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:
DPAsyncMPClientañadecurrent_wave/lb_engines/eep_scaling_cachey_ensure_stats_update_task;DPLBAsyncMPClientsobrescribecall_utility_asyncpara hacer LB interno(DP init:1200-1240). La clase base no carga esta lógica de DP. - El hook
output_handlerse 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_refantes de invocarlo y no capturarselfen 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.