Skip to content

RayDistributedExecutor: mit Ray Worker über mehrere Knoten und GPUs verteilen

源码版本v0.25.1

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 PlacementGroupSchedulingStrategy an 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 ein collective_rpc("adjust_rank", args=(rerank_mapping,)) die Ranks gerade.
  • Ray Compiled Graph für den Forward:Statt pro Schritt execute_model per RPC aufzurufen, wird ein forward_dag aufgebaut(ray_executor.py:527-620); bei PP=2/TP=4 werden vier Worker zu InputNode -> TP group -> MultiOutputNode verkettet, und Zwischentensoren wandern per with_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_variables gebündelt an jeden Worker gepusht(ray_executor.py:305-327) — ausgenommen WORKER_SPECIFIC_ENV_VARS, damit sich die Worker nicht gegenseitig verschmutzen.
  • execute_model in zwei Phasen:Bei aktivem Sampler führt execute_model den Forward nicht sofort aus, sondern speichert nur self.scheduler_output; erst wenn sample_tokens aufgerufen 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-97initialize_ray_cluster, Placement Group abholen, _init_workers_ray aufrufen, has_connector, uses_sampler und den Platzhalter scheduler_output initialisieren.
  • shutdown:99-113 — nach forward_dag.teardown() folgt ein schrittweises ray.kill(worker); das Log weist darauf hin, dass SIGTERM erwartet wird.
  • _init_workers_ray:143-213 — pro bundle_indices ein ray.remote(...)(RayWorkerWrapper).remote(rpc_rank=rank); GPUs mit num_gpus=num_gpus, andere Plattformen über resources={ray_device_key: num_gpus}.
  • rerank:217-258 — über ray.get die Knoten-IP jedes Actors abholen, sortieren nach „gleicher Driver-Knoten zuerst → weniger Actors pro Knoten zuerst → kleinere IP zuerst“ und anschließend per collective_rpc("adjust_rank", ...) den rpc_rank korrigieren.
  • init_worker + load_model:343-368all_kwargs zusammenstellen und gebündelt collective_rpc("init_worker") + collective_rpc("init_device") + collective_rpc("load_model").
  • pp_tp_workers:370-378 — füllt self.pp_tp_workers in einer PP×TP-Matrix, Grundlage für den späteren Compiled DAG.
  • execute_model + sample_tokens:390-432 — Zweiphasen: execute_model speichert nur scheduler_output, der echte Lauf findet in sample_tokens über _execute_dag statt.
  • _execute_dag:434-468 — beim ersten Aufruf wird self.forward_dag = self._compiled_ray_dag(...) erzeugt; forward_dag.execute((scheduler_output, grammar_output)) liefert refs zurück — ohne PP blockiert refs[0].get() bis das Ergebnis da ist, nichtblockierend dagegen liefert es FutureWrapper.
  • collective_rpc:470-495 — serialisiert die Methode (cloudpickle.dumps), ruft pro Worker execute_method.remote(...) auf und sammelt mit ray.get(ray_worker_outputs, timeout=timeout) ein; non_block=True liefert FutureWrapper.
  • _compiled_ray_dag:527-620 — baut aus InputNode + worker.execute_model_ray.bind(...) einen DAG, mit with_tensor_transport(transport) zwischen PP-Stufen, und schließt mit experimental_compile(...) ab.
  • check_health:625-628 — in der aktuellen Version direkt return; ein TODO ist für eine echte Ray-Actor-Gesundheitsprüfung hinterlegt.
  • RayWorkerWrapper:56-91 — erbt von WorkerWrapperBase, ergänzt adjust_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:

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

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:

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

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_ray prüft explizit n_nodes == n_ips und wirft bei Verletzung einen RuntimeError, der auf VLLM_HOST_IP verweist(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 falsches get_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_model feststellt, dass self.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_installation verlangt 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_executor erzwingt bei TPU/XPU VLLM_USE_RAY_COMPILED_DAG_CHANNEL_TYPE=shm(ray_executor.py:73-75), um NVIDIA NCCL auszulassen.
  • check_health ist aktuell leer:RayDistributedExecutor.check_health ist direkt return(ray_executor.py:625-628) — stürzt ein Actor wirklich ab, wird die Exception über ray.get propagiert; aktives Heartbeat gibt es derzeit nicht.
  • Shutdown: harte Kill-Methode:Nach forward_dag.teardown() folgt ray.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.

Siehe offizielle Dokumentation: vLLM 文档 · README.