RayDistributedExecutor: mit Ray Worker über mehrere Knoten und GPUs verteilen
Verantwortung
Im Single-GPU-Fall ist UniProcExecutor mit einem einzigen Prozess erledigt — sobald jedoch TP>1 oder ein knotenübergreifendes Setup ins Spiel kommt, wird Inter-Prozess-Kommunikation unabdingbar. RayDistributedExecutor(ray_executor.py:64-68) ist die Implementierung für das Ray-Backend: Sie wickelt jeden Worker in einen Ray-Actor ein, verteilt ihn über die durch die Placement Group vorgegebenen GPUs und baut mit dem Ray Compiled Graph eine Pipeline, die den Forward zusammenführt. uses_ray = True, supports_pp = True(ray_executor.py:67-68) — weist auf Ray-Orchestrierung hin und signalisiert Unterstützung für Pipeline-Parallelität (pipeline parallelism).
Wie UniProcExecutor erbt auch diese Klasse von Executor und lässt die Signaturen der Basisklasse für execute_model/sample_tokens/profile/sleep unverändert; überschrieben werden nur _init_executor, collective_rpc, execute_model, check_health und shutdown(ray_executor.py:70-97). Der Unterschied liegt einzig im RPC-Pfad: Im selben Prozess ist es ein direkter Methodenaufruf, in Ray ist es worker.execute_method.remote(...)(ray_executor.py:484-489).
Entwurfsmotivation
- Ray-Actor entspricht einem Worker:Auf jeder GPU läuft ein
RayWorkerWrapper-Actor(ray_executor.py:200-213); Ray übernimmt Scheduling, Fehlerbehandlung und Ressourcenbuchhaltung — der Executor spawnt selbst keine Prozesse. - Placement Group fixiert die Topologie:Die Actors werden mit
PlacementGroupSchedulingStrategyan Bundles gebunden(ray_executor.py:192-196) — die GPU-Affinität steht vor dem Start fest und wandert während der Laufzeit nicht. - Bevorzugte Sortierung gleicher Knoten:Die Reihenfolge der Actor-Erzeugung wird von Ray bestimmt und kann ungeordnet sein; daher werden nach der Erzeugung alle Actors nach IP neu sortiert, wobei Actors auf dem Knoten des Drivers zuerst kommen(
ray_executor.py:234-258) — anschließend setzt eincollective_rpc("adjust_rank", args=(rerank_mapping,))die Ranks gerade. - Ray Compiled Graph für den Forward:Statt pro Schritt
execute_modelper RPC aufzurufen, wird einforward_dagaufgebaut(ray_executor.py:527-620); bei PP=2/TP=4 werden vier Worker zuInputNode -> TP group -> MultiOutputNodeverkettet, und Zwischentensoren wandern perwith_tensor_transport(transport="nccl"|"shm")über NCCL oder Shared Memory(ray_executor.py:582-591). - Umgebungsvariablen automatisch kopieren:vLLM-bezogene Umgebungsvariablen des Driver-Prozesses werden per
update_environment_variablesgebündelt an jeden Worker gepusht(ray_executor.py:305-327) — ausgenommenWORKER_SPECIFIC_ENV_VARS, damit sich die Worker nicht gegenseitig verschmutzen. - execute_model in zwei Phasen:Bei aktivem Sampler führt
execute_modelden Forward nicht sofort aus, sondern speichert nurself.scheduler_output; erst wennsample_tokensaufgerufen wird, folgt_execute_dag(ray_executor.py:390-432) — so lassen sich die Grammar-Bitmask-Berechnung des Schedulers/Samplers und der Forward überlappen.
Schlüsseldateien
RayDistributedExecutor._init_executor:64-97—initialize_ray_cluster, Placement Group abholen,_init_workers_rayaufrufen,has_connector,uses_samplerund den Platzhalterscheduler_outputinitialisieren.shutdown:99-113— nachforward_dag.teardown()folgt ein schrittweisesray.kill(worker); das Log weist darauf hin, dass SIGTERM erwartet wird._init_workers_ray:143-213— probundle_indiceseinray.remote(...)(RayWorkerWrapper).remote(rpc_rank=rank); GPUs mitnum_gpus=num_gpus, andere Plattformen überresources={ray_device_key: num_gpus}.rerank:217-258— überray.getdie Knoten-IP jedes Actors abholen, sortieren nach „gleicher Driver-Knoten zuerst → weniger Actors pro Knoten zuerst → kleinere IP zuerst“ und anschließend percollective_rpc("adjust_rank", ...)den rpc_rank korrigieren.init_worker + load_model:343-368—all_kwargszusammenstellen und gebündeltcollective_rpc("init_worker")+collective_rpc("init_device")+collective_rpc("load_model").pp_tp_workers:370-378— fülltself.pp_tp_workersin einer PP×TP-Matrix, Grundlage für den späteren Compiled DAG.execute_model + sample_tokens:390-432— Zweiphasen:execute_modelspeichert nurscheduler_output, der echte Lauf findet insample_tokensüber_execute_dagstatt._execute_dag:434-468— beim ersten Aufruf wirdself.forward_dag = self._compiled_ray_dag(...)erzeugt;forward_dag.execute((scheduler_output, grammar_output))liefert refs zurück — ohne PP blockiertrefs[0].get()bis das Ergebnis da ist, nichtblockierend dagegen liefert esFutureWrapper.collective_rpc:470-495— serialisiert die Methode (cloudpickle.dumps), ruft pro Workerexecute_method.remote(...)auf und sammelt mitray.get(ray_worker_outputs, timeout=timeout)ein;non_block=TrueliefertFutureWrapper._compiled_ray_dag:527-620— baut ausInputNode+worker.execute_model_ray.bind(...)einen DAG, mitwith_tensor_transport(transport)zwischen PP-Stufen, und schließt mitexperimental_compile(...)ab.check_health:625-628— in der aktuellen Version direktreturn; ein TODO ist für eine echte Ray-Actor-Gesundheitsprüfung hinterlegt.RayWorkerWrapper:56-91— erbt vonWorkerWrapperBase, ergänztadjust_rank,get_node_ip,get_node_and_physical_gpu_ids,execute_method(verpackt Ausnahmen in einem „möglicher Deadlock“-Log und wirft sie dann weiter).
Datenfluss
Während des Starts durchläuft _init_executor den Pfad _init_workers_ray — der längste Abschnitt: Bundle wählen, RayWorkerWrapper-Actors anhand der Placement Group erzeugen, dann über ray.get die Knoten-IP jedes Actors abholen für Sortierung und 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))Zur Laufzeit, wenn EngineCore.step bei execute_model ankommt und uses_sampler=True sowie ein nicht-leerer Batch vorliegt, speichert der Executor nur den scheduler_output und kehrt zurück(ray_executor.py:401-407). Erst der Aufruf von sample_tokens führt wirklich _execute_dag aus — beim ersten Mal wird der Compiled DAG erzeugt:
# 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])
...Die Control-Plane-RPCs (initialize_from_config, compile_or_warm_up_model, determine_available_memory usw.) laufen alle über collective_rpc — also pro Worker ein execute_method.remote(...)(ray_executor.py:484-489) und danach ray.get zum Einsammeln. Das ist ein vom Datenpfad (Forward über Compiled DAG) komplett getrennter Kanal.
Grenzen und Fehler
- Knoten-IPs müssen eindeutig sein:
_init_workers_rayprüft explizitn_nodes == n_ipsund wirft bei Verletzung einenRuntimeError, der aufVLLM_HOST_IPverweist(ray_executor.py:295-303). - Loopback-Adresse auf einem einzelnen Knoten:Bei mehreren Actors auf demselben Knoten wird
driver_ip = "127.0.0.1"gesetzt, damit ein falschesget_ip()über die falsche Netzwerkkarte die Kommunikation zwischen Workern nicht scheitern lässt(ray_executor.py:329-341). - Falsche Reihenfolge der Zustandsmaschine löst Fehler aus:Wenn
execute_modelfeststellt, dassself.scheduler_output is not None, wirft es sofort „sample_tokens() must be called after execute_model() returns None“(ray_executor.py:395-399). - Ray Compiled Graph: Version und Abhängigkeiten:
_check_ray_cgraph_installationverlangt ray>=2.43.0,ray[cgraph]und optional cupy(ray_executor.py:497-525) — sonst scheitert der Start bereits mit einem Fail-Fast. - TPU/XPU: Shared Memory:
_init_executorerzwingt bei TPU/XPUVLLM_USE_RAY_COMPILED_DAG_CHANNEL_TYPE=shm(ray_executor.py:73-75), um NVIDIA NCCL auszulassen. - check_health ist aktuell leer:
RayDistributedExecutor.check_healthist direktreturn(ray_executor.py:625-628) — stürzt ein Actor wirklich ab, wird die Exception überray.getpropagiert; aktives Heartbeat gibt es derzeit nicht. - Shutdown: harte Kill-Methode:Nach
forward_dag.teardown()folgtray.kill(worker)(ray_executor.py:111-112) — nicht gespeicherte Zustände im Worker gehen verloren; in der Dokumentation wird explizit darauf hingewiesen, dass SIGTERM-Logs erwartet werden.
Zusammenfassung
RayDistributedExecutor ist die Executor-Implementierung von v1 für Multi-Node- und Multi-GPU-Szenarien. Sie wickelt jeden Worker in einen Ray-Actor ein, fixiert die Topologie über die Placement Group und baut mit dem Ray Compiled Graph einen Forward-DAG, der TP/PP zusammenführt; die Control-Plane-RPCs laufen weiterhin über den execute_method-Pfad des Actors. Gegenüber UniProcExecutor bleibt die an EngineCore exponierte Schnittstelle identisch — nur der RPC-Mechanismus unter der Haube ist ein anderer. Wie der Executor abstrahiert ist, steht in /executor/executor; wie der Worker intern den Forward ausführt, in /worker/worker und /worker/gpu-model-runner.