Skip to content

EngineCoreProc : enveloppe de processus arrière ZMQ

源码版本v0.25.1

Responsabilités

EngineCoreProc est une sous-classe d'EngineCore (core.py:896-897) qui fait tourner tout son contenu dans un processus arrière séparé (background process). Le docstring est direct : ZMQ-wrapper for running EngineCore in background process.. Il fait trois choses : au démarrage, poignée de main avec le frontal pour récupérer l'adresse des sockets ZMQ ; lance deux threads d'IO pour relier socket et deux queue.Queue (core.py:915-916) ; puis run_busy_loop tourne sur le thread principal et boucle sur : prendre une requête depuis input_queuestep_fn() pour ordonnancer et lancer le forward → pousser les outputs dans output_queue (core.py:1259-1267).

Cette couche est la clé de l'architecture multi-processus v1. Le processus frontal (SyncMPClient ou AsyncMPClient) sérialise les EngineCoreRequest via ZMQ ; le thread input côté EngineCoreProc les reçoit et fait put_nowait dans input_queue, et la boucle principale, entre deux ordonnancements, les retire et les passe à _handle_client_request pour dispatch (core.py:1372-1405). Les outputs circulent en sens inverse : les EngineCoreOutputs produits par step_fn() sont poussés dans output_queue, le thread output les retire, les sérialise et les renvoie au frontal via ZMQ. Le socket IO ZMQ libère le GIL à l'envoi comme à la réception, donc ces deux threads d'IO se chevauchent vraiment avec le forward GPU (core.py:974-1001), et la sérialisation/désérialisation se trouve masquée derrière le forward.

Motivation de conception

  • Isolation de processus : EngineCoreProc tourne dans un processus séparé, donc un crash Python frontal n'emporte pas les workers déjà engagés dans le GPU, et un OOM worker ne tue pas le processus d'API. C'est la raison centrale pour laquelle v1 extirpe EngineCore : en v0, LLMEngine faisait tout dans le même processus, un crash à un endroit faisait tout tomber.
  • Double buffer ZMQ + queue.Queue : le socket ZMQ n'est pas utilisé directement par la boucle principale ; on intercale une queue.Queue (core.py:915-916). Le socket IO vit dans le thread d'IO, la boucle principale ne fait que queue.get(block=...), sans être bloquée par un recv socket ni par un verrou interne ZMQ ; et quand le socket IO libère le GIL, la boucle principale peut le récupérer pour lancer le forward.
  • Deux threads d'IO : process_input_sockets et process_output_sockets (core.py:980-1001) sont indépendants : le thread input gère la désérialisation + put_nowait, le thread output gère get + sérialisation + send_multipart, sans se bloquer l'un l'autre.
  • Poignée de main au démarrage (handshake) : dans le constructeur, _perform_handshakes (core.py:932-938) récupère un EngineZmqAddresses, puis le thread input envoie une EngineCoreReadyResponse sur le ROUTER/DEALER pour signaler au frontal « je suis prêt, voilà max_model_len, num_gpu_blocks, kv_cache_size_tokens » (core.py:1525-1546). Côté frontal, MPClient.__init__ attend ce message ready ; s'il n'arrive pas, le démarrage est considéré comme échoué.
  • Point d'entrée run_engine_core : cette fonction est la cible de multiprocessing.Process(target=run_engine_core, ...) (core.py:1154-1224). Elle instancie EngineCoreProc, enregistre SIGTERM/SIGINT, puis appelle run_busy_loop. Toute exception passe par except Exception_send_engine_deadraise, et à la sortie le finally appelle engine_core.shutdown() (core.py:1226-1242).

Fichiers clés

  • EngineCoreProc class:896-910class EngineCoreProc(EngineCore): + constante de classe ENGINE_CORE_DEAD = b"ENGINE_CORE_DEAD".
  • EngineCoreProc.__init__:903-1009 — crée input_queue/output_queue, tensor IPC receiver, _perform_handshakes, super().__init__, lance les deux threads d'IO.
  • queues + executor_fail_callback:915-919input_queue = queue.Queue[...], output_queue = queue.Queue[...], le callback d'échec d'executor pousse EXECUTOR_FAILED dans input_queue.
  • IO threads startup:974-1001 — lancent les deux threads daemon process_input_sockets et process_output_sockets ; ready_event n'est set qu'après que le thread input est prêt.
  • run_engine_core:1154-1224 — point d'entrée du processus, choisit EngineCoreProc ou DPEngineCoreProc, enregistre les handlers de signaux, appelle run_busy_loop.
  • signal handlers:1204-1222wakeup_engine pousse WAKEUP dans input_queue, SIGTERM/SIGINT passe shutdown_state à REQUESTED.
  • run_busy_loop:1259-1267while self._handle_shutdown(): _process_input_queue(); _process_engine_step() puis raise SystemExit à la sortie du shutdown.
  • _process_input_queue:1269-1298input_queue.get(block=True) quand il n'y a rien à faire, drain non-bloquant sinon.
  • _process_engine_step:1300-1317 — appelle step_fn(), output_queue.put_nowait(output), post_step ; si aucun forward mais scheduler encore actif, time.sleep(0.001) rend le GIL.
  • _handle_client_request:1372-1405 — dispatch par EngineCoreRequestType : WAKEUP/ADD/ABORT/UTILITY/EXECUTOR_FAILED.
  • _send_engine_dead:1470-1482 — pousse ENGINE_CORE_DEAD dans output_queue et attend que le thread output l'envoie avant de sortir.
  • process_input_sockets:1484-1587 — corps du thread input, utilise zmq.Poller sur DEALER + XSUB coord optionnel, désérialise puis put_nowait dans input_queue.
  • process_output_sockets:1589-1654 — corps du thread output, prend les outputs via output_queue.get(), encoder.encode_into + send_multipart(copy=False, track=True) pour renvoyer au frontal en zero-copy.

Flux de données

La requête entre par le socket ZMQ input, est désérialisée par le thread input et poussée dans input_queue ; la boucle principale la retire dans _process_input_queue et la passe à _handle_client_request pour dispatch :

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))

Notez que ABORT est poussé à la fois dans aborts_queue et input_queue, parce que l'abort est idempotent côté scheduler (core.py:1276-1279) ; cela permet à la boucle principale de drainer aborts_queue entre deux steps, pour éviter de fuite des requêtes déjà abortées.

Une fois la requête récupérée par la boucle principale, _handle_client_request dispatche selon 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 est un canal RPC générique : le frontal envoie le nom de méthode + args msgspec, EngineCoreProc l'invoque via getattr par réflexion, et le résultat repart par output_queue. _invoke_utility_method gère aussi le cas où la valeur renvoyée est un Future (core.py:1440-1446) ; l'output n'est renvoyé qu'une fois l'utilitaire asynchrone terminé.

Limites et échecs

  • Notification de crash du processus : run_engine_core, dans son except Exception, appelle engine_core._send_engine_dead() (core.py:1229-1235) qui pousse la bytestring ENGINE_CORE_DEAD dans output_queue (core.py:1470-1474) ; le thread output, en voyant cette valeur spéciale, la diffuse à tous les sockets d'output (core.py:1625-1628) ; le MPClient frontal lève alors une erreur au lieu d'attendre sans fin.
  • Callback d'échec d'Executor : __init__ enregistre la closure executor_fail_callback auprès de model_executor (core.py:917-919) ; en cas d'erreur côté worker, le callback pousse EXECUTOR_FAILED dans input_queue, et la boucle principale, dans _handle_client_request, lève raise RuntimeError("Executor failed.") (core.py:1400-1401) ; le processus tout entier est ensuite attrapé par le except de run_engine_core qui déclenche _send_engine_dead.
  • Mode shutdown : _handle_shutdown consulte shutdown_state (core.py:1324-1360) ; à l'état REQUESTED, shutdown_timeout==0 déclenche le mode abort (tue immédiatement toutes les requêtes en vol), non nul déclenche le mode drain (attend la fin de toutes les requêtes). Les ADD et UTILITY qui arrivent pendant le shutdown sont explicitement rejetés (core.py:1407-1432) : _reject_add_in_shutdown renvoie un output d'abort au client, _reject_utility_in_shutdown renvoie un UtilityOutput avec failure_message="Server shutting down".
  • Mort du thread input pendant le démarrage : __init__ boucle sur ready_event.wait(timeout=10) dans l'attente de la readiness du thread input (core.py:1003-1009) ; si au bout de 10 secondes le thread n'est toujours pas prêt et qu'il est mort, on raise RuntimeError("Input socket thread died during startup") ; s'il est encore vivant mais pas prêt, on continue d'attendre le message READY du DP Coordinator.
  • Envoi zero-copy + MessageTracker : le thread output utilise send_multipart(buffers, copy=False, track=True) (core.py:1646-1654) qui renvoie un MessageTracker ; les trackers + références outputs + buffers non encore envoyés sont accrochés à une deque pending, pour empêcher le GC de récupérer les backing buffers avant leur envoi. reuse_buffers limite le nombre de buffers réutilisés à len(sockets) + 1 pour éviter une croissance mémoire non bornée.
  • WAKEUP est un no-op : le handler de signaux pousse WAKEUP dans input_queue via wakeup_engine, uniquement pour réveiller la boucle principale bloquée dans input_queue.get(block=True) (core.py:1377-1378) ; en soi il ne fait rien, le vrai traitement du shutdown se fait dans _handle_shutdown en lisant shutdown_state.
  • Restauration des signaux dans le finally : run_engine_core, dans son finally, fait signal.signal(signal.SIGTERM, signal.SIG_DFL) (core.py:1236-1242) pour remettre le handler par défaut, puis appelle engine_core.shutdown() pour libérer model_executor / scheduler / état distribué. shutdown appelle aussi gc.unfreeze() (core.py:652-656) pour rendre au GC le startup heap figé par gc.freeze() au démarrage, sinon le mode in-process (tests unitaires) fuit de la mémoire GPU.

Résumé

EngineCoreProc est le mur porteur de l'architecture multi-processus v1 : il hérite d'EngineCore, tourne dans un processus séparé avec ZMQ + deux queue.Queue + deux threads d'IO, de façon à ce que le socket IO et le forward GPU se chevauchent vraiment ; run_engine_core est le point d'entrée du processus, chargé de l'instantiation, de l'enregistrement des signaux, du filet de sécurité en cas d'exception et de shutdown. Toute communication entre fronts (SyncMPClient / AsyncMPClient) et back (EngineCore) passe par cette couche. Pour aller plus loin : /engine/engine-core pour voir comment le cœur du moteur lui-même est assemblé, /engine/llm-engine pour la coquille frontale, les trois implémentations client à /client/inproc-mp, et l'intérieur du scheduler à /scheduler/scheduler.

Voir la documentation officielle : Documentation vLLM · README