Skip to content

EngineCoreClient: dos caminos de cliente, en proceso y multiproceso

源码版本v0.25.1

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 de MPClient.__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 recoge BackgroundResources.__call__, sin fugas.
  • Las interfaces síncrona y asíncrona comparten el mismo protocolo ZMQ: SyncMPClient usa una queue.Queue para mover los frames del output socket al hilo principal, y AsyncMPClient usa una asyncio.Queue, pero MPClient comparte 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.sleep L324-L328), porque en modo en-proceso no hay un planificador concurrente al que «esperar a que quede inactivo para dormir»; solo queda abort. 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_rpc o save_sharded_state, no pueden ser fire-and-forget como add_request. call_utility genera un call_id, mete un future en self.utility_results, y cuando llega un UtilityOutput desde el output socket, _process_utility_output saca el future con el call_id y lo rellena con set_result. El mecanismo es idéntico en los dos clientes; solo cambia el tipo de future (concurrent.futures.Future frente a asyncio.Future).
  • Propagación de engine dead: cuando el subproceso muere de forma inesperada, el hilo monitor pone BackgroundResources.engine_dead a True, y cualquier _send_input posterior es bloqueado por ensure_alive y lanza EngineDeadError. En el lado del output socket, si llega un único frame ENGINE_CORE_DEAD, validate_alive también levanta el mismo flag, de modo que tanto el lado de envío como el de recepción lo perciben.

Archivos clave

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:

python
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=False se 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__ usa try/finally + un flag success; si falla, invoca self._finalizer() para disparar BackgroundResources.__call__, que detiene los subprocesos ya arrancados por launch_core_engines(__init__ finally:499-649).
  • InprocClient no soporta la pausa wait(InprocClient.sleep:324-328). En modo en-proceso nadie espera a estar inactivo, así que solo se puede abort la 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 con validate_alive al recibir el frame único ENGINE_CORE_DEAD(validate_alive:454-457). Ambos levantan BackgroundResources.engine_dead.
  • Tolerancia a InvalidStateError en el future utility: _process_utility_output captura asyncio.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 a close_sockets antes 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 hace self.resources.output_socket = None(socket ownership transfer:846-847), para que BackgroundResources y 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 ajuste VLLM_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.

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