Skip to content

EngineCoreProc: wrapper de proceso en segundo plano con ZMQ

源码版本v0.25.1

Responsabilidades

EngineCoreProc es una subclase de EngineCore (core.py:896-897), y ejecuta todo dentro de un proceso en segundo plano (background process) independiente. Su docstring lo deja claro: ZMQ-wrapper for running EngineCore in background process.. Hace tres cosas: durante el arranque del proceso hace handshake con el frontend para obtener las direcciones de los sockets ZMQ, levanta dos hilos de IO que puentean los sockets con dos queue.Queue (core.py:915-916), y luego run_busy_loop en el hilo principal va sacando requests de input_queue repetidamente → ejecuta step_fn() con scheduling y forward → mete los outputs en output_queue (core.py:1259-1267).

Esta capa es clave en la arquitectura multiproceso de v1. El proceso frontend (SyncMPClient o AsyncMPClient) envía vía ZMQ un EngineCoreRequest serializado; el hilo de input de EngineCoreProc lo recibe y hace put_nowait en input_queue, y el loop principal lo saca en un hueco del scheduling y lo despacha con _handle_client_request (core.py:1372-1405). Los outputs fluyen en sentido inverso: los EngineCoreOutputs que produce step_fn() se dejan en output_queue, el hilo de output los saca, los serializa y los envía de vuelta al frontend por ZMQ. El IO de los sockets ZMQ libera el GIL al enviar/recibir, así que esos dos hilos de IO pueden solaparse de verdad con el forward en GPU (core.py:974-1001); la serialización/deserialización también queda escondida detrás del forward.

Motivación de diseño

  • Aislamiento de procesos: EngineCoreProc corre en un proceso independiente; si el Python del frontend se cae no arrastra a los workers que ya están corriendo forward, y un OOM en un worker no se lleva por delante al proceso de API. Esa es la razón central por la que v1 saca EngineCore afuera: en v0 LLMEngine hacía todo en el mismo proceso y un fallo en cualquier parte tumbaba todo.
  • Doble buffer con ZMQ + queue.Queue: el socket ZMQ no lo usa el loop principal directamente, hay una capa intermedia con queue.Queue (core.py:915-916). El IO del socket vive en hilos de IO, el loop principal solo hace queue.get(block=...), evitando que el loop principal se quede bloqueado en un recv del socket o arrastrado por un lock interno de ZMQ; al mismo tiempo, cuando el IO del socket libera el GIL, el loop principal puede retomarlo de inmediato para correr forward.
  • Dos hilos de IO: process_input_sockets y process_output_sockets (core.py:980-1001) son independientes; el hilo de input se encarga de deserializar + put_nowait, el de output hace get + serializar + send_multipart, y no se bloquean entre sí.
  • Handshake de arranque: en el constructor, _perform_handshakes (core.py:932-938) obtiene un EngineZmqAddresses, y luego el hilo de input envía un EngineCoreReadyResponse por ROUTER/DEALER para avisarle al frontend "ya estoy arriba, este es el max_model_len, num_gpu_blocks, kv_cache_size_tokens" (core.py:1525-1546). El frontend espera este mensaje ready en MPClient.__init__; si no llega, lo trata como fallo de arranque.
  • run_engine_core como entrada del proceso: esta función es el target de multiprocessing.Process(target=run_engine_core, ...) (core.py:1154-1224). Instancia EngineCoreProc, registra SIGTERM/SIGINT, y al final llama run_busy_loop. Cualquier excepción cae en except Exception_send_engine_deadraise; cuando el proceso sale, en finally se llama engine_core.shutdown() (core.py:1226-1242).

Archivos clave

  • EngineCoreProc class:896-910 — definición de class EngineCoreProc(EngineCore): + la constante de clase ENGINE_CORE_DEAD = b"ENGINE_CORE_DEAD".
  • EngineCoreProc.__init__:903-1009 — levanta input_queue/output_queue, el receptor de tensor IPC, _perform_handshakes, super().__init__, y arranca los dos hilos de IO.
  • queues + executor_fail_callback:915-919input_queue = queue.Queue[...], output_queue = queue.Queue[...]; el callback de fallo del executor mete EXECUTOR_FAILED en input_queue.
  • IO threads startup:974-1001 — arranca los dos daemon threads process_input_sockets y process_output_sockets; ready_event solo se setea cuando el hilo de input queda ready.
  • run_engine_core:1154-1224 — entrada del proceso, elige EngineCoreProc o DPEngineCoreProc, registra signal handlers, y llama run_busy_loop.
  • signal handlers:1204-1222wakeup_engine mete WAKEUP en input_queue; SIGTERM/SIGINT cambian shutdown_state a REQUESTED.
  • run_busy_loop:1259-1267while self._handle_shutdown(): _process_input_queue(); _process_engine_step(), y al terminar el shutdown lanza raise SystemExit.
  • _process_input_queue:1269-1298 — sin trabajo bloquea en input_queue.get(block=True), con trabajo hace un drain no bloqueante.
  • _process_engine_step:1300-1317 — llama step_fn(), output_queue.put_nowait(output), post_step; si no hubo forward pero el scheduler aún tiene trabajo, hace time.sleep(0.001) para ceder el GIL.
  • _handle_client_request:1372-1405 — despacha EngineCoreRequestType: WAKEUP/ADD/ABORT/UTILITY/EXECUTOR_FAILED.
  • _send_engine_dead:1470-1482 — mete ENGINE_CORE_DEAD en output_queue y espera a que el hilo de output lo envíe antes de salir.
  • process_input_sockets:1484-1587 — cuerpo del hilo de input, usa zmq.Poller sobre DEALER + un XSUB de coord opcional, deserializa y hace put_nowait en input_queue.
  • process_output_sockets:1589-1654 — cuerpo del hilo de output, saca outputs con output_queue.get(), hace encoder.encode_into + send_multipart(copy=False, track=True) para enviarlos al frontend sin copiar.

Flujo de datos

Las requests entran por el socket de input de ZMQ, el hilo de input las deserializa y las mete en input_queue; el loop principal las saca en _process_input_queue y se las pasa a _handle_client_request para despacharlas:

python
# vllm/v1/engine/core.py L1556-L1587
while True:
    for input_socket, _ in poller.poll():
        # (RequestType, RequestData)
        type_frame, *data_frames = input_socket.recv_multipart(copy=False)
        # NOTE(yongji): ignore READY message sent by DP coordinator
        # that is used to notify newly started engines
        if type_frame.buffer == b"READY":
            assert input_socket == coord_socket
            continue
        request_type = EngineCoreRequestType(bytes(type_frame.buffer))

        # Deserialize the request data.
        request: Any
        if request_type == EngineCoreRequestType.ADD:
            req: EngineCoreRequest = add_request_decoder.decode(data_frames)
            try:
                request = self.preprocess_add_request(req)
            except Exception:
                self._handle_request_preproc_error(req)
                continue
        else:
            request = generic_decoder.decode(data_frames)

            if request_type == EngineCoreRequestType.ABORT:
                # Aborts are added to *both* queues, allows us to eagerly
                # process aborts while also ensuring ordering in the input
                # queue to avoid leaking requests. This is ok because
                # aborting in the scheduler is idempotent.
                self.aborts_queue.put_nowait(request)

        # Push to input queue for core busy loop.
        self.input_queue.put_nowait((request_type, request))

Notar que ABORT se mete tanto en aborts_queue como en input_queue, porque abort es idempotente en el scheduler (core.py:1276-1279); así el loop principal puede también hacer drain de aborts_queue entre steps, evitando fugar requests ya abortados.

Cuando el loop principal recibe la request, _handle_client_request la despacha según EngineCoreRequestType (core.py:1372-1405):

python
# vllm/v1/engine/core.py L1377-L1401
if request_type == EngineCoreRequestType.WAKEUP:
    return
elif request_type == EngineCoreRequestType.ADD:
    req, request_wave = request
    if self._reject_add_in_shutdown(req):
        return
    self.add_request(req, request_wave)
elif request_type == EngineCoreRequestType.ABORT:
    self.abort_requests(request)
elif request_type == EngineCoreRequestType.UTILITY:
    client_idx, call_id, method_name, args = request
    if self._reject_utility_in_shutdown(client_idx, call_id, method_name):
        return
    output = UtilityOutput(call_id)
    # Lazily look-up utility method so that failure will be handled/returned.
    get_result = lambda: (
        (method := getattr(self, method_name))
        and method(*self._convert_msgspec_args(method, args))
    )
    enqueue_output = lambda out: self.output_queue.put_nowait(
        (client_idx, EngineCoreOutputs(utility_output=out))
    )
    self._invoke_utility_method(method_name, get_result, output, enqueue_output)
elif request_type == EngineCoreRequestType.EXECUTOR_FAILED:
    raise RuntimeError("Executor failed.")

UTILITY es un canal RPC genérico: el frontend envía el nombre del método + args msgspec, EngineCoreProc lo invoca por reflexión con getattr, y el resultado vuelve por output_queue. _invoke_utility_method también maneja el caso en el que se devuelve un Future (core.py:1440-1446): el output solo se envía cuando la utilidad async termina.

Límites y fallos

  • Notificación de caída del proceso: run_engine_core en el except Exception llama engine_core._send_engine_dead() (core.py:1229-1235), que mete en output_queue la cadena de bytes ENGINE_CORE_DEAD (core.py:1470-1474); cuando el hilo de output la recibe, la difunde a todos los output sockets (core.py:1625-1628), y el MPClient del frontend al recibirla lanza un error en vez de esperar indefinidamente.
  • Callback de fallo del executor: en __init__ se registra el closure executor_fail_callback en model_executor (core.py:917-919); cuando un worker falla internamente, el callback mete un EXECUTOR_FAILED en input_queue, y el loop principal en _handle_client_request lo recibe y lanza raise RuntimeError("Executor failed.") (core.py:1400-1401); el proceso entero termina siendo capturado por el except de run_engine_core y cae en _send_engine_dead.
  • Modo shutdown: _handle_shutdown mira shutdown_state (core.py:1324-1360); en estado REQUESTED, si shutdown_timeout==0 va por abort y mata de inmediato todas las requests en vuelo, si no es 0 va por drain y espera a que terminen todas. Durante el shutdown las nuevas ADD y UTILITY se rechazan explícitamente (core.py:1407-1432): _reject_add_in_shutdown devuelve un output de abort al cliente, y _reject_utility_in_shutdown devuelve un UtilityOutput con failure_message="Server shutting down".
  • Muerte del hilo de input: en __init__ el loop ready_event.wait(timeout=10) espera a que el hilo de input quede ready (core.py:1003-1009); si pasan 10 segundos sin ready y el hilo está muerto, lanza raise RuntimeError("Input socket thread died during startup"); si solo no está ready pero sigue vivo, sigue esperando el mensaje READY del DP Coordinator.
  • Envío zero-copy + MessageTracker: el hilo de output usa send_multipart(buffers, copy=False, track=True) (core.py:1646-1654), que devuelve un MessageTracker; los trackers que no terminaron de enviar, junto con referencias a outputs y buffers, se cuelgan en un deque pending, para que el GC no reclame el backing buffer que aún no se envió. reuse_buffers acota a len(sockets) + 1 la cantidad de buffers reutilizables, evitando que la memoria crezca sin límite.
  • WAKEUP es no-op: el signal handler mete WAKEUP en input_queue a través de wakeup_engine, solo para despertar al loop principal bloqueado en input_queue.get(block=True) (core.py:1377-1378); no hace nada por sí mismo, el manejo real del shutdown lo hace _handle_shutdown mirando shutdown_state.
  • finally restaura señales: run_engine_core en finally ejecuta signal.signal(signal.SIGTERM, signal.SIG_DFL) (core.py:1236-1242), devuelve los signal handlers a su valor por defecto, y luego llama engine_core.shutdown() para liberar model_executor / scheduler / estado distribuido. shutdown también llama gc.unfreeze() (core.py:652-656), para devolver al GC el startup heap que se había congelado con gc.freeze() en el arranque; sin esto, el modo en el mismo proceso (unit test) filtra VRAM.

Resumen

EngineCoreProc es el muro de carga de la arquitectura multiproceso de v1: hereda de EngineCore, corre en un proceso independiente con ZMQ + dos queue.Queue + dos hilos de IO, y deja que el IO del socket y el forward en GPU se solapen de verdad; run_engine_core es la entrada del proceso, encargada de instanciar, registrar señales, atrapar excepciones y llamar shutdown. Toda la comunicación entre el frontend (SyncMPClient / AsyncMPClient) y el backend (EngineCore) pasa por esta capa. Para ver cómo se ensambla el núcleo del motor, continúa en /engine/engine-core; para ver la shell fina del frontend, en /engine/llm-engine; las tres implementaciones de cliente están en /client/inproc-mp, y el interior del scheduler en /scheduler/scheduler.

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