EngineCoreClient: dos caminos de cliente, en proceso y multiproceso
Responsabilidades
El frontend de vLLM (LLM / AsyncLLM / las entradas de serving) no toca EngineCore directamente, sino a través de una capa de cliente que empuja las peticiones y extrae las salidas. Esa capa es EngineCoreClient, y unifica dos dimensiones ortogonales —«síncrono o asíncrono», «mismo proceso o subproceso»— bajo el mismo conjunto de métodos abstractos. Todas las implementaciones viven en vllm/v1/engine/core_client.py.
Las dos implementaciones clave son InprocClient y MPClient. La primera construye EngineCore directamente en el proceso actual, lo invoca paso a paso, sin busy loop y sin ZMQ; funciona sobre todo como una capa fina para el estilo V0 LLMEngine.add_request()/step(). La segunda envuelve EngineCore como subproceso EngineCoreProc y se comunica por un par de sockets ZMQ ROUTER/PULL para enviar peticiones y recibir salidas; el hilo del frontend solo habla con sockets. MPClient es la clase base, con dos materializaciones: SyncMPClient (síncrono, arranca un hilo en segundo plano para leer del output socket) y AsyncMPClient (asyncio, lanza una tarea para el mismo socket), esta última tratada en otra página.
EngineCoreClient.make_client es una fábrica estática que elige la implementación según dos booleanos, multiprocess_mode y asyncio_mode. La combinación (asyncio=True, multiprocess=False) lanza directamente NotImplementedError: EngineCore no es async-friendly por sí mismo, y el modo asíncrono exige multiproceso.
Motivación de diseño
- Desacoplar el ciclo de vida del frontend y de EngineCore: llamar a EngineCore en el mismo proceso es simple, pero para correr varios EngineCore en paralelo (paralelismo de datos) o para soportar EP elástico (añadir o quitar ranks en línea) hay que empujar EngineCore a un proceso aparte y que el frontend solo retenga sockets. La línea
weakref.finalize(self, self.resources)dentro deMPClient.__init__es la clave: incluso si la construcción lanza una excepción a mitad, el subproceso en segundo plano y el context ZMQ los recogeBackgroundResources.__call__, sin fugas. - Las interfaces síncrona y asíncrona comparten el mismo protocolo ZMQ:
SyncMPClientusa unaqueue.Queuepara mover los frames del output socket al hilo principal, yAsyncMPClientusa unaasyncio.Queue, peroMPClientcomparte encoder/decoder, handshake de ready, engine monitor y el registro de futures utility. Las diferencias se relegan a las dos subclases. InprocClient.sleep(mode="wait")se prohíbe directamente(InprocClient.sleepL324-L328), porque en modo en-proceso no hay un planificador concurrente al que «esperar a que quede inactivo para dormir»; solo quedaabort. Esta frontera sí está soportada en modo subproceso: el subproceso tiene su propio busy loop y su máquina de estados de planificación.- Las llamadas utility van por un registro de futures: los RPC que deben devolver un valor, como
collective_rpcosave_sharded_state, no pueden ser fire-and-forget comoadd_request.call_utilitygenera uncall_id, mete un future enself.utility_results, y cuando llega unUtilityOutputdesde el output socket,_process_utility_outputsaca el future con elcall_idy lo rellena conset_result. El mecanismo es idéntico en los dos clientes; solo cambia el tipo de future (concurrent.futures.Futurefrente aasyncio.Future). - Propagación de engine dead: cuando el subproceso muere de forma inesperada, el hilo monitor pone
BackgroundResources.engine_deada True, y cualquier_send_inputposterior es bloqueado porensure_alivey lanzaEngineDeadError. En el lado del output socket, si llega un único frameENGINE_CORE_DEAD,validate_alivetambién levanta el mismo flag, de modo que tanto el lado de envío como el de recepción lo perciben.
Archivos clave
EngineCoreClient + make_client:71-105— clase abstracta y fábrica estática; eligeInprocClient/SyncMPClient/AsyncMPClientsegúnmultiprocess_mode/asyncio_modemake_async_mp_client:107-132— fábrica asíncrona; segúndata_parallel_sizey si hay LB externo, eligeAsyncMPClient/DPAsyncMPClient/DPLBAsyncMPClientInprocClient.__init__ / get_output:276-292— construyeEngineCoredirectamente;get_outputes una llamada síncrona astep_fn()+post_step()InprocClient.add_request / abort / shutdown:297-306— reenvío directo aself.engine_core, sin serialización ni socketInprocClient.sleep rechaza wait:324-328— el modo en-proceso no soporta pausawait, soloabortBackgroundResources:370-458— dataclass con__call__como callback deweakref.finalize; cierra sockets, detiene tasks y apaga el engine managervalidate_alive:454-457— al recibir el frame únicoENGINE_CORE_DEAD, marcaengine_deady lanzaEngineDeadErrorMPClient.__init__:480-649— crea context/sockets ZMQ, opcionalmentelaunch_core_engines, hace handshake de ready y arranca el monitor; revierte si fallaMPClient.shutdown:651-662— tras el finalizerdetach(), detiene el engine manager y ejecuta la limpieza deBackgroundResourcesstart_engine_core_monitor:685-712— hilo demoniomonitor_engine_liveness; al morir el engine, poneengine_deady disparaself.shutdown()_apply_ready_response:714-754— decodificaEngineCoreReadyResponsey rellenamax_model_len/num_gpu_blocks/block_sizedesde el subproceso hacia elvllm_configdel frontendSyncMPClient.__init__:779-847— crea laqueue.Queue, arranca el hilo demonioprocess_outputs_socketque encuesta el output socket y el shutdown socketSyncMPClient.get_output / _send_input / call_utility:849-881— bloquea conoutputs_queue.get(); al enviar, usatrackpara obtenerMessageTrackery evitar la recolección prematura de tensores
Flujo de datos
SyncMPClient saca salidas mediante un hilo demonio que hace poll sobre el output_socket de ZMQ, decodifica los frames y los empuja a una queue.Queue. El hilo del frontend llama a get_output(), que no es más que queue.get(): una adaptación del socket no bloqueante de ZMQ a una interfaz síncrona bloqueante de Python. Este process_outputs_socket es el corazón de la ruta síncrona:
def process_outputs_socket():
assert isinstance(out_socket, zmq.Socket)
shutdown_socket = ctx.socket(zmq.PAIR)
try:
shutdown_socket.bind(shutdown_path)
poller = zmq.Poller()
poller.register(shutdown_socket, zmq.POLLIN)
poller.register(out_socket, zmq.POLLIN)
while True:
socks = poller.poll()
if not socks:
continue
if len(socks) == 2 or socks[0][0] == shutdown_socket:
# shutdown signal, exit thread.
break
frames = out_socket.recv_multipart(copy=False)
resources.validate_alive(frames)
outputs: EngineCoreOutputs = decoder.decode(frames)
if outputs.utility_output:
_process_utility_output(outputs.utility_output, utility_results)
else:
outputs_queue.put_nowait(outputs)
except Exception as e:
outputs_queue.put_nowait(e)
finally:
# Close sockets.
shutdown_socket.close(linger=0)
out_socket.close(linger=0)(SyncMPClient.process_outputs_socket:808-836). Hay tres detalles: primero, el shutdown se hace con un zmq.PAIR inproc aparte atado a shutdown_path; BackgroundResources.__call__ envía un byte vacío a ese socket al cerrar para notificar al hilo, y este no se quede esperando para siempre en out_socket. Segundo, validate_alive convierte un único frame ENGINE_CORE_DEAD en EngineDeadError. Tercero, la respuesta utility (p. ej. el retorno de collective_rpc) y las salidas normales viajan en el mismo frame; las primeras se inyectan en el future de utility_results[call_id], las segundas entran en la cola.
El camino de envío es más simple: SyncMPClient._send_input envía directamente los frames (identity, request_type, *encoder.encode(request)) al ROUTER; cuando hay tensores, se usa track=True para obtener un MessageTracker y meterlo en pending_messages, evitando que Python recoja el tensor subyacente antes de tiempo(_send_input:861-873).
El handshake se hace en MPClient.__init__(esperando ready:615-634): para cada engine_rank se genera una identidad de 2 bytes little-endian y se hace poll sobre el input_socket a la espera del EngineCoreReadyResponse del subproceso, con timeout controlado por VLLM_ENGINE_READY_TIMEOUT_S. _apply_ready_response rellena en el vllm_config del frontend los num_gpu_blocks, el block_size ya alineado y el max_model_len autoajustado que el subproceso realmente ha reservado: esa es la capacidad real que el frontend puede usar para planificar y limitar.
Límites y fallos
asyncio_mode=True, multiprocess_mode=Falsese rechaza directamente (make_client:91-95). EngineCore no es async-friendly internamente y esa combinación no está implementada.- Si la construcción falla, hay que recoger el subproceso:
MPClient.__init__usatry/finally+ un flagsuccess; si falla, invocaself._finalizer()para dispararBackgroundResources.__call__, que detiene los subprocesos ya arrancados porlaunch_core_engines(__init__ finally:499-649). InprocClientno soporta la pausawait(InprocClient.sleep:324-328). En modo en-proceso nadie espera a estar inactivo, así que solo se puedeabortla petición actual antes de dormir.- La muerte del subproceso engine core se percibe por ambos lados: el lado de envío con
ensure_alive(ensure_alive:670-672), el de recepción convalidate_aliveal recibir el frame únicoENGINE_CORE_DEAD(validate_alive:454-457). Ambos levantanBackgroundResources.engine_dead. - Tolerancia a InvalidStateError en el future utility:
_process_utility_outputcapturaasyncio.InvalidStateError(_process_utility_output:769-776), para el caso en el que la tarea llamadora se ha cancelado y el future ya está cancelled. - El cierre del output socket debe preceder a
context.term:BackgroundResources.__call__llama explícitamente aclose_socketsantes de seguir(sync case:437-452); si no, la terminación del context ZMQ se cuelga. - En
SyncMPClient, el output_socket lo cierra el hilo: al final del constructor se haceself.resources.output_socket = None(socket ownership transfer:846-847), para queBackgroundResourcesy el hilo demonio no intenten cerrarlo a la vez. - El timeout del handshake de ready da un error legible:
Timed out waiting for engine core processes to start. This is often caused by slow weight loading for large models.(ready timeout:622-630) — se le indica al usuario que ajusteVLLM_ENGINE_READY_TIMEOUT_S.
Resumen
InprocClient y MPClient parecen muy distintos por dentro, pero hacia fuera se presentan ambos como EngineCoreClient —un add_request y un get_output, sin que el frontend se preocete por si EngineCore es un objeto en proceso o un subproceso con ZMQ. El modo en-proceso es lo más simple posible: una referencia a objeto. El modo subproceso encapsula el ROUTER/PULL de ZMQ, el handshake, el monitor, la tabla de futures y la propiedad del socket. La variante asíncrona AsyncMPClient se desglosa en /client/async-mp; EngineCore en sí está en /engine/engine-core, y el busy loop y la señal de shutdown de EngineCoreProc en /engine/engine-core-proc.