Skip to content

EngineCoreClient : deux voies client, in-process et multi-processus

源码版本v0.25.1

Responsabilités

Le frontal de vLLM (LLM / AsyncLLM / entrée serving) ne touche pas directement à EngineCore : il passe par une couche client pour pousser les requêtes et récupérer les sorties. Cette couche, c'est EngineCoreClient, qui unifie sur un même ensemble de méthodes abstraites les deux dimensions orthogonales « synchrone ou asynchrone » et « même processus ou sous-processus ». Le fichier vllm/v1/engine/core_client.py contient toutes les implémentations.

Les deux implémentations clés sont InprocClient et MPClient. La première construit EngineCore directement dans le processus courant, appelle un step à la fois, sans busy loop ni ZMQ — principalement une coquille fine pour LLMEngine.add_request()/step() à la V0. La seconde enveloppe EngineCore dans un sous-processus EngineCoreProc et utilise une paire de sockets ZMQ ROUTER/PULL pour émettre les requêtes et recevoir les sorties ; le thread appelant côté frontal ne parle qu'aux sockets. MPClient est la classe de base, concrétisée par SyncMPClient (synchrone, lance un thread en arrière-plan pour consommer le socket d'output) et AsyncMPClient (asyncio, lance une task en arrière-plan sur le même socket), cette dernière étant traitée dans une page dédiée.

EngineCoreClient.make_client est une fabrique statique qui choisit l'implémentation à partir de deux booléens, multiprocess_mode et asyncio_mode. La combinaison (asyncio=True, multiprocess=False) lève directement NotImplementedError : EngineCore lui-même n'est pas asyncio-friendly, l'asynchrone ne fonctionne qu'avec le mode multi-processus.

Motivation de conception

  • Découpler la vie du frontal et d'EngineCore : appeler EngineCore dans le même processus est simple, mais pour en faire tourner plusieurs en parallèle (parallélisme de données) ou pour EP élastique (ajout/retrait de rank en ligne), il faut pousser EngineCore dans un processus séparé et que le frontal ne tienne que des sockets. La ligne weakref.finalize(self, self.resources) dans MPClient.__init__ est centrale : même si une exception est levée en cours de construction, le sous-processus en arrière-plan et le contexte ZMQ sont récupérés par BackgroundResources.__call__, sans fuite.
  • Interface synchrone et asynchrone partagent le même protocole ZMQ : SyncMPClient utilise queue.Queue pour remonter les frames du socket d'output vers le thread principal, AsyncMPClient utilise asyncio.Queue, mais MPClient mutualise encoder/decoder, handshake ready, engine monitor et registre de utility futures. Les différences sont repoussées dans les deux sous-classes.
  • sleep(mode="wait") est interdit pour InprocClient (InprocClient.sleep L324-L328) : en mode in-process il n'y a pas d'ordonnanceur concurrent qui puisse « attendre l'inactivité avant de dormir », seul abort fonctionne. Cette frontière est en revanche supportée en mode sous-processus, qui a son propre busy loop et sa machine à états d'ordonnancement.
  • Les appels utility passent par un registre de futures : les RPC qui doivent renvoyer une valeur (collective_rpc, save_sharded_state, etc.) ne peuvent pas être fire-and-forget comme add_request. call_utility génère un call_id, range la future dans self.utility_results, et quand un UtilityOutput revient par le socket d'output, _process_utility_output récupère la future via call_id et la complète avec set_result. Ce mécanisme est identique sur les clients sync et async, seule la nature de la future diffère (concurrent.futures.Future vs asyncio.Future).
  • Propagation de engine dead : quand le sous-processus meurt inopinément, le thread monitor positionne BackgroundResources.engine_dead à True, et tout _send_input ultérieur est bloqué par ensure_alive qui lève EngineDeadError. Côté socket d'output, si une frame unique ENGINE_CORE_DEAD arrive, validate_alive lève le même drapeau, garantissant que l'émission et la réception voient toutes deux l'état mort.

Fichiers clés

Flux de données

SyncMPClient tire ses sorties via un thread daemon qui poll le output_socket ZMQ, décode les frames et les pousse dans une queue.Queue. Le thread frontal appelle get_output(), qui est un simple queue.get() — l'adaptation du socket ZMQ non bloquant en interface synchrone bloquante Python. Ce process_outputs_socket est le cœur de la voie synchrone :

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). Trois détails à noter. D'abord, le shutdown passe par un socket inproc zmq.PAIR séparé bindé sur shutdown_path ; à la fermeture, BackgroundResources.__call__ envoie un octet vide sur ce socket pour notifier le thread de sortir, afin d'éviter qu'il reste bloqué sur out_socket. Ensuite, validate_alive convertit une frame unique ENGINE_CORE_DEAD en EngineDeadError. Enfin, les paquets utility (retour de collective_rpc, etc.) et les outputs ordinaires passent par la même frame : les premiers sont distribués directement dans la future utility_results[call_id], les seconds vont dans la queue.

L'émission de requête est plus simple : SyncMPClient._send_input envoie directement les frames multiples (identity, request_type, *encoder.encode(request)) sur le ROUTER, et en présence d'un tensor buffer utilise track=True pour obtenir un MessageTracker, raccroché à pending_messages pour empêcher le GC Python de libérer le tensor sous-jacent (_send_input:861-873).

Le handshake se fait dans MPClient.__init__ (attente ready:615-634) : pour chaque engine_rank on génère une identity little-endian de 2 octets, on poll sur input_socket en attendant le EngineCoreReadyResponse du sous-processus, avec un timeout piloté par VLLM_ENGINE_READY_TIMEOUT_S. _apply_ready_response remonte dans le vllm_config frontal le num_gpu_blocks calculé par le sous-processus, le block_size aligné, et le max_model_len auto-ajusté — c'est la véritable capacité exploitable pour l'ordonnancement et la limitation côté frontal.

Limites et échecs

  • asyncio_mode=True, multiprocess_mode=False est rejeté (make_client:91-95). EngineCore n'est pas async-friendly en interne, cette combinaison n'est pas implémentée.
  • Échec à la construction doit recycler le sous-processus : MPClient.__init__ utilise try/finally + un drapeau success ; en cas d'échec il invoque self._finalizer() qui déclenche BackgroundResources.__call__, stoppant les sous-processus lancés par launch_core_engines (__init__ finally:499-649).
  • InprocClient ne supporte pas la pause wait (InprocClient.sleep:324-328). En mode in-process personne n'attend l'inactivité, il n'y a que abort des requêtes en cours avant sleep.
  • La mort du sous-processus engine core est perçue des deux côtés : côté émission via ensure_alive (ensure_alive:670-672), côté réception via validate_alive qui reçoit une frame ENGINE_CORE_DEAD (validate_alive:454-457). Les deux chemins lèvent BackgroundResources.engine_dead.
  • Tolérance à InvalidStateError sur les futures utility : _process_utility_output capture asyncio.InvalidStateError (_process_utility_output:769-776), pour couvrir le cas où la task appelante a été annulée et a déjà annulé la future.
  • Fermeture du socket d'output avant context.term : BackgroundResources.__call__ fait explicitement close_sockets avant la suite (sync case:437-452), sinon la terminaison du contexte ZMQ se bloque.
  • La fermeture du output_socket de SyncMPClient est assumée par le thread : à la fin de la construction, self.resources.output_socket = None (socket ownership transféré:846-847), pour éviter que BackgroundResources et le thread daemon ne le ferment simultanément.
  • Message lisible en cas de timeout du handshake ready : Timed out waiting for engine core processes to start. This is often caused by slow weight loading for large models. (ready timeout:622-630) — indique directement à l'utilisateur d'ajuster VLLM_ENGINE_READY_TIMEOUT_S.

Résumé

Les deux voies InprocClient et MPClient paraissent très différentes mais exposent toutes deux la surface EngineCoreClient — un add_request, un get_output, et le code frontal n'a pas à savoir si EngineCore est un objet dans le même processus ou un sous-processus derrière ZMQ. Le mode in-process est minimal à l'extrême : une simple référence d'objet ; le mode sous-processus encapsule ZMQ ROUTER/PULL, le handshake, le monitor, la table de futures et la propriété des sockets. La variante asynchrone AsyncMPClient est détaillée sur /client/async-mp ; EngineCore lui-même est sur /engine/engine-core, et le busy loop / shutdown d'EngineCoreProc sur /engine/engine-core-proc.

Voir la documentation officielle : Documentation vLLM · README