RayDistributedExecutor : utiliser Ray pour répartir les workers sur plusieurs nœuds et GPU
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
RayWorkerWrapperpar 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 appellecollective_rpc("adjust_rank", args=(rerank_mapping,))pour rectifier les ranks. - Ray Compiled Graph pour la passe avant : au lieu d'invoquer
execute_modelpar RPC à chaque étape, on construit uneforward_dag(ray_executor.py:527-620). Pour PP=2/TP=4, on enchaîne les quatre workers enInputNode -> TP group -> MultiOutputNode, les tenseurs intermédiaires transitant parwith_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 deWORKER_SPECIFIC_ENV_VARSafin d'éviter la pollution croisée entre workers. execute_modelscindé en deux étapes : si le sampler est activé,execute_modelne s'exécute pas immédiatement ; il stocke seulementself.scheduler_outputet attend l'appelsample_tokenspour 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-97—initialize_ray_cluster, récupération du placement group, appel de_init_workers_ray, initialisation dehas_connector,uses_sampleret emplacementscheduler_output.shutdown:99-113— aprèsforward_dag.teardown(), appelleray.kill(worker)un par un ; le journal indique que SIGTERM est un comportement attendu._init_workers_ray:143-213— pour chaquebundle_indices,ray.remote(...)(RayWorkerWrapper).remote(rpc_rank=rank); pour les GPUnum_gpus=num_gpus, sinonresources={ray_device_key: num_gpus}.rerank:217-258—ray.getré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 », puiscollective_rpc("adjust_rank", ...)rectifie le rpc_rank.init_worker + load_model:343-368— une foisall_kwargsassemblé, enchaînecollective_rpc("init_worker")+collective_rpc("init_device")+collective_rpc("load_model").pp_tp_workers:370-378— remplitself.pp_tp_workersen deux dimensions PP×TP en vue de la construction du DAG compilé.execute_model + sample_tokens:390-432— two-stage :execute_modelstockescheduler_output, et la véritable exécution se fait danssample_tokensvia_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, renvoieFutureWrapper.collective_rpc:470-495— sérialise la méthode (cloudpickle.dumps), appelleexecute_method.remote(...)sur chaque worker, puisray.get(ray_worker_outputs, timeout=timeout); avecnon_block=TruerenvoieFutureWrapper._compiled_ray_dag:527-620— compose le DAG avecInputNode+worker.execute_model_ray.bind(...), utilisewith_tensor_transport(transport)pour transmettre les tenseurs entre étages PP, puisexperimental_compile(...).check_health:625-628— dans la version actuelle, fait directementreturn; un TODO est prévu pour une véritable vérification de santé des actors Ray.RayWorkerWrapper:56-91— hérite deWorkerWrapperBase, ajouteadjust_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 :
# 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 :
# 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_rayvérifie explicitementn_nodes == n_ips, sinon lève uneRuntimeErroret suggère d'inspecterVLLM_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 queget_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_modellève immédiatement « sample_tokens() must be called after execute_model() returns None » siself.scheduler_output is not None(ray_executor.py:395-399). - Version et dépendances de Ray Compiled Graph :
_check_ray_cgraph_installationexige 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_executordétecte TPU/XPU, il forceVLLM_USE_RAY_COMPILED_DAG_CHANNEL_TYPE=shm(ray_executor.py:73-75) pour éviter NCCL de NVIDIA. - check_health est actuellement vide :
RayDistributedExecutor.check_healthse contente dereturn(ray_executor.py:625-628). En cas de plantage d'un actor, la remontée dépend des exceptions levées parray.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