Skip to content

RayDistributedExecutor : utiliser Ray pour répartir les workers sur plusieurs nœuds et GPU

源码版本v0.25.1

Responsabilités

Sur un seul GPU, UniProcExecutor s'en sort avec un seul processus. Mais dès que TP>1 ou que l'on franchit un nœud, il faut de la communication inter-processus. RayDistributedExecutor(ray_executor.py:64-68) est l'implémentation du backend Ray : elle emballe chaque worker dans un Ray actor, les répartit sur les GPU indiqués par le placement group, puis utilise Ray Compiled Graph pour construire un pipeline qui enchaîne la passe avant. uses_ray = True, supports_pp = True(ray_executor.py:67-68) indique d'une part qu'il s'agit de l'orchestration via Ray, d'autre part la prise en charge du parallélisme de pipeline (pipeline parallelism).

Comme UniProcExecutor, elle hérite d'Executor et ne modifie pas les signatures de execute_model/sample_tokens/profile/sleep de la classe de base ; seuls _init_executor, collective_rpc, execute_model, check_health et shutdown sont redéfinis(ray_executor.py:70-97). La différence se situe uniquement dans l'acheminement des RPC : en intra-processus on appelle directement la méthode, en Ray on passe par worker.execute_method.remote(...)(ray_executor.py:484-489).

Motivation de conception

  • Ray actor en correspondance un-à-un avec chaque worker : un actor RayWorkerWrapper par GPU(ray_executor.py:200-213). Ray prend en charge la planification, la tolérance aux pannes et le accounting de ressources ; l'exécuteur ne spawn pas lui-même de processus.
  • Placement group verrouille la topologie : les actors sont liés à un bundle via PlacementGroupSchedulingStrategy(ray_executor.py:192-196). L'affinité GPU est fixée avant le démarrage et ne dérive pas à l'exécution.
  • Tri priorisant le même nœud : l'ordre de création des actors est décidé par Ray (potentiellement non déterministe). Une fois tous les actors créés, on les réordonne par IP, en plaçant en premier ceux qui partagent le nœud du driver(ray_executor.py:234-258), puis on appelle collective_rpc("adjust_rank", args=(rerank_mapping,)) pour rectifier les ranks.
  • Ray Compiled Graph pour la passe avant : au lieu d'invoquer execute_model par RPC à chaque étape, on construit une forward_dag(ray_executor.py:527-620). Pour PP=2/TP=4, on enchaîne les quatre workers en InputNode -> TP group -> MultiOutputNode, les tenseurs intermédiaires transitant par with_tensor_transport(transport="nccl"|"shm") via NCCL ou la mémoire partagée(ray_executor.py:582-591).
  • Réplication automatique des variables d'environnement : les variables d'environnement vLLM du driver sont poussées en lot vers chaque worker via update_environment_variables(ray_executor.py:305-327), à l'exclusion de WORKER_SPECIFIC_ENV_VARS afin d'éviter la pollution croisée entre workers.
  • execute_model scindé en deux étapes : si le sampler est activé, execute_model ne s'exécute pas immédiatement ; il stocke seulement self.scheduler_output et attend l'appel sample_tokens pour déclencher _execute_dag(ray_executor.py:390-432). Cela permet au scheduler et au calcul du bitmask de grammar du sampler de se chevaucher avec la passe avant.

Fichiers clés

  • RayDistributedExecutor._init_executor:64-97initialize_ray_cluster, récupération du placement group, appel de _init_workers_ray, initialisation de has_connector, uses_sampler et emplacement scheduler_output.
  • shutdown:99-113 — après forward_dag.teardown(), appelle ray.kill(worker) un par un ; le journal indique que SIGTERM est un comportement attendu.
  • _init_workers_ray:143-213 — pour chaque bundle_indices, ray.remote(...)(RayWorkerWrapper).remote(rpc_rank=rank) ; pour les GPU num_gpus=num_gpus, sinon resources={ray_device_key: num_gpus}.
  • rerank:217-258ray.get récupère l'IP de chaque actor, puis on trie selon « même nœud que le driver d'abord → moins d'actors sur le nœud d'abord → plus petite IP d'abord », puis collective_rpc("adjust_rank", ...) rectifie le rpc_rank.
  • init_worker + load_model:343-368 — une fois all_kwargs assemblé, enchaîne collective_rpc("init_worker") + collective_rpc("init_device") + collective_rpc("load_model").
  • pp_tp_workers:370-378 — remplit self.pp_tp_workers en deux dimensions PP×TP en vue de la construction du DAG compilé.
  • execute_model + sample_tokens:390-432 — two-stage : execute_model stocke scheduler_output, et la véritable exécution se fait dans sample_tokens via _execute_dag.
  • _execute_dag:434-468 — au premier appel, self.forward_dag = self._compiled_ray_dag(...) ; forward_dag.execute((scheduler_output, grammar_output)) renvoie des refs. Si PP est désactivé, refs[0].get() bloque pour récupérer le résultat ; en mode non bloquant, renvoie FutureWrapper.
  • collective_rpc:470-495 — sérialise la méthode (cloudpickle.dumps), appelle execute_method.remote(...) sur chaque worker, puis ray.get(ray_worker_outputs, timeout=timeout) ; avec non_block=True renvoie FutureWrapper.
  • _compiled_ray_dag:527-620 — compose le DAG avec InputNode + worker.execute_model_ray.bind(...), utilise with_tensor_transport(transport) pour transmettre les tenseurs entre étages PP, puis experimental_compile(...).
  • check_health:625-628 — dans la version actuelle, fait directement return ; un TODO est prévu pour une véritable vérification de santé des actors Ray.
  • RayWorkerWrapper:56-91 — hérite de WorkerWrapperBase, ajoute adjust_rank, get_node_ip, get_node_and_physical_gpu_ids, execute_method (enveloppe l'exception dans un log « possible deadlock » avant de la relancer).

Flux de données

Au démarrage, _init_executor passe par _init_workers_ray, la phase la plus longue : sélection des bundles, création des actors RayWorkerWrapper selon le placement group, puis ray.get pour récupérer l'IP de chaque actor afin de les trier et de procéder au rerank :

python
# vllm/v1/executor/ray_executor.py L191-L213
for rank, bundle_id in enumerate(bundle_indices):
    scheduling_strategy = PlacementGroupSchedulingStrategy(
        placement_group=placement_group,
        placement_group_capture_child_tasks=True,
        placement_group_bundle_index=bundle_id,
    )

    if current_platform.ray_device_key == "GPU":
        worker = ray.remote(
            num_cpus=0,
            num_gpus=num_gpus,
            scheduling_strategy=scheduling_strategy,
            **ray_remote_kwargs,
        )(RayWorkerWrapper).remote(rpc_rank=rank)
    else:
        worker = ray.remote(
            num_cpus=0,
            num_gpus=0,
            resources={current_platform.ray_device_key: num_gpus},
            scheduling_strategy=scheduling_strategy,
            **ray_remote_kwargs,
        )(RayWorkerWrapper).remote(rpc_rank=rank)

    worker_metadata.append(RayWorkerMetaData(worker=worker, created_rank=rank))

À l'exécution, quand EngineCore.step appelle execute_model, si uses_sampler=True et qu'il y a des tokens dans le batch, l'exécuteur se contente de stocker scheduler_output et de renvoyer(ray_executor.py:401-407) ; ce n'est qu'à l'appel de sample_tokens que _execute_dag s'exécute réellement, en construisant le DAG compilé au premier appel :

python
# vllm/v1/executor/ray_executor.py L434-L452
def _execute_dag(self, scheduler_output, grammar_output, non_block=False):
    if self.forward_dag is None:
        self.forward_dag = self._compiled_ray_dag(enable_asyncio=False)

    refs = self.forward_dag.execute((scheduler_output, grammar_output))

    if not self.has_connector:
        if not non_block:
            output = refs[0].get()
            detach_zero_copy_from_model_runner_output(output)
            return output
        return FutureWrapper(refs[0])
    ...

Les RPC du plan de contrôle (initialize_from_config, compile_or_warm_up_model, determine_available_memory, etc.) passent tous par collective_rpc, c'est-à-dire execute_method.remote(...) sur chaque worker(ray_executor.py:484-489), puis récupérés via ray.get. Cela constitue un chemin distinct du plan de données (passe avant via le DAG compilé).

Limites et échecs

  • Les IP de nœud doivent être uniques : _init_workers_ray vérifie explicitement n_nodes == n_ips, sinon lève une RuntimeError et suggère d'inspecter VLLM_HOST_IP(ray_executor.py:295-303).
  • Bouclage sur adresse de loopback en mono-nœud : lorsque plusieurs actors cohabitent sur le même nœud, driver_ip = "127.0.0.1", afin d'éviter que get_ip() ne sélectionne une mauvaise carte réseau et ne provoque des échecs de communication entre workers(ray_executor.py:329-341).
  • Erreur si l'ordre de la machine à états est violé : execute_model lève immédiatement « sample_tokens() must be called after execute_model() returns None » si self.scheduler_output is not None(ray_executor.py:395-399).
  • Version et dépendances de Ray Compiled Graph : _check_ray_cgraph_installation exige ray>=2.43.0, ray[cgraph], et optionnellement cupy(ray_executor.py:497-525), sinon échec rapide au démarrage.
  • TPU/XPU en mémoire partagée : quand _init_executor détecte TPU/XPU, il force VLLM_USE_RAY_COMPILED_DAG_CHANNEL_TYPE=shm(ray_executor.py:73-75) pour éviter NCCL de NVIDIA.
  • check_health est actuellement vide : RayDistributedExecutor.check_health se contente de return(ray_executor.py:625-628). En cas de plantage d'un actor, la remontée dépend des exceptions levées par ray.get ; il n'y a pas de heartbeat actif.
  • shutdown brutale : après forward_dag.teardown(), ray.kill(worker)(ray_executor.py:111-112). Tout état non flushé dans le worker est perdu ; la documentation précise que le log SIGTERM est attendu.

Résumé

RayDistributedExecutor est l'implémentation de l'exécuteur (executor) v1 pour les scénarios multi-nœuds et multi-GPU. Elle emballe chaque worker dans un Ray actor, s'appuie sur le placement group pour verrouiller la topologie, et construit un DAG de passe avant via Ray Compiled Graph pour enchaîner TP/PP ; les RPC du plan de contrôle continuent de passer par execute_method des actors. L'interface exposée à EngineCore est identique à celle d'UniProcExecutor, seul le mécanisme RPC sous-jacent change. Pour la façon dont l'exécuteur est abstrait, voir /executor/executor ; pour la passe avant interne du worker, voir /worker/worker et /worker/gpu-model-runner.

Voir la documentation officielle : Documentation vLLM · README