Skip to content

[RFC]: Revamp Ray Distributed Executor Backend (from Ray team) #35848

Description

@jeffreywang88

Summary

Replace the current RayExecutor distributed executor backend which uses Compiled Graph with a new RayExecutorV2 that launches workers as plain Ray remote actors and follows the MultiprocExecutor communication model: ZMQ/SHM-backed MessageQueue for the control plane and NCCL collective communication for the data plane.

Motivation

1. Compiled Graph’s optimizations are no longer the primary bottleneck

The Compiled Graph backend was introduced to minimize scheduling overhead by compiling the execution DAG in advance and to provide NCCL abstractions for Ray actor communication. Since then, two developments have removed the need for both:

  • Async scheduling (v0.15.0) overlaps control plane scheduling with GPU execution, eliminating the per-step scheduling latency that Compiled Graph was designed to hide.
  • torch.distributed already provides NCCL collectives natively. vLLM's MultiprocExecutor uses these directly for worker-to-worker communication.

Compiled Graph does offer additional optimizations, notably overlapping communication with computation, that are not yet implemented in vLLM. However, these can be achieved independently through torch.distributed APIs and CUDA stream pipelining without requiring compiled graph.

The remaining benefits of Compiled Graph do not justify the stability and maintenance costs described below.

2. Maintenance burden

The Compiled Graph backend has unresolved stability issues, primarily NCCL hangs and out-of-order delivery, that have persisted across multiple Ray releases. Moreover, the Ray Core team has limited resources to maintain Compiled Graph.

3. Divergence from the MultiprocExecutor golden path

MultiprocExecutor is vLLM's primary executor and has first-class compatibility with new features (e.g. async scheduling). The existing Ray executor requires separate implementation work for each feature due to its fundamentally different communication model. This creates ongoing friction for the community: contributors must verify and adapt features for a second backend with different semantics.

Ideally, the Ray executor should share the same control plane as MultiprocExecutor so that features work on both backends without additional effort.

4. Multi-node support & granular placement for collocation and resource-aware scheduling

Ray's core strength is cluster orchestration: process placement, fault detection, and resource management across nodes. Rather than deprecating Ray from vLLM entirely, we propose replacing the Compiled Graph backend with a simpler design that uses Ray solely as a process launcher and placement manager. The control plane (MessageQueue) and data plane (torch.distributed / NCCL) remain identical to MultiprocExecutor, while Ray enables simpler cross-node parallelism, granular placement of GPU workers for collocation and resource-aware scheduling that MultiprocExecutor cannot support. For example, SkyRL utilizes VLLM_BUNDLE_INDICES to schedule inference engine actor (vLLM) with a specific placement.

Best reasons NOT to remove Ray from vLLM entirely

  1. Simpler cross-node parallelism setup: In theory, vLLM could be set up with driver / headless worker pattern in a Ray cluster, but that's more involved that directly spinning vLLM with Ray as the distributed executor backend.
  2. Process cleanup: When using Ray as the external orchestrator with vLLM's mp backend, there will be zombie vLLM processes upon fault / OOM scenarios.
  3. RL applications -- enable Ray RDT to orchestrate inference / trainer workers: Ray's RDT assumes that GPU workers are part of the transfer and their state and placements are managed by Ray's GCS. RDT is incompatible with compiled graph, making it difficult integrate vLLM with RL application while using Ray as the orchestrator.

Proposed Change

1. Goals / Non-goals

1.1 Goal

  • Align Ray executor with MultiprocExecutor by sharing the same control and data plane communication
  • Multi-node pipeline parallelism: Enable cross-node PP/TP by leveraging Ray’s cluster orchestration.
  • Placement group flexibility: Accept client-provided placement groups and support granular bundle index placement for deployment platforms that manage GPU placement externally.

1.2 Non-goals

  • Replacing Ray entirely: Ray remains vLLM’s solution for multi-node orchestration and granular placement.
  • Communication-computation overlap: Porting Compiled Graph’s ability to overlap NCCL communication with GPU computation is out of scope for this RFC.

2. Proposed changes

2.1 API

[Selected] Option 1: Inherit from MultiprocExecutor
Implement RayExecutorV2 as a child class of MultiprocExecutor to reuse maximal code.

  • Pros: Zero code duplication, best long-term maintainability, and zero risk to the existing MultiprocExecutor.
  • Cons: Any change to MultiprocExecutor could break RayExecutorV2 unexpectedly. We can mitigate this with comprehensive test coverage, and one of the goals of this RFC is to align RayExecutorV2 closely with MultiprocExecutor.
  class RayExecutorV2(MultiprocExecutor):                                                                                                                                                                                                                                                                                                                   
      def _init_executor(self) -> None:
          """
          Stage 1: initialize_ray_cluster() → get/create PlacementGroup
          Stage 2: query PG table, sort bundles, assign ranks + local_ranks
          Stage 3: create broadcast MessageQueue
          Stage 4: spawn RayWorkerProc actors into PG bundles
          Stage 5: ray.get(init_refs) → collect response MQ handles
          Stage 6: wait_until_ready() barrier
          Stage 7: actor.run.remote() → store run_refs, start health monitor
          """

      def start_worker_monitor(self):
          """
          ray.wait() on the ObjectRef to run() which runs the worker busy loop.
          """

      def shutdown(self) -> None:
          """ray.kill(actor) for each actor. Set  shutdown_event."""

Option 2: Independent RayExecutorV2
Implement an isolated RayExecutorV2 although there are substantial common pieces between RayExecutorV2 and MultiprocExecutor, e.g. execute_model, collective_rpc, etc.

  • Pros: Zero risk to existing MultiprocExecutor.
  • Cons: Bug fixes or features added to one executor must be re-applied to the other.

Option 3: Extract common pieces to a new class MessageQueueExecutor
Refactor the duplicated logic in MultiprocExecutor and RayExecutorV2 into a unified MessageQueueExecutor.

  • Pros: Reduces code duplication and improves long-term maintainability.
  • Cons: Requires modifying vLLM core components, which may introduce stability risks.

2.2 Placement group support

Placement group creation and bundle assignment must be done prior to actor launch because ranks are required during actor initialization to create NCCL collective groups and collective groups are immutable after creation.

  1. Get or create the placement group: Reuse initialize_ray_cluster().
    a. If the client passes one, use it directly.
    b. If we’re running inside a placement group, use it.
    c. Create placement group with PACK strategy with bundle 0 pinned to the driver node with label selector.
  2. Sort bundles: Driver node gets the lowest rank, and same node bundles get contiguous ranks.
  3. Compute global and per-node local rank.
  4. Classify local vs. remote for MessageQueue creation (more details in the following section).
  5. Spawn actors with final ranks. No adjust_rank RPC needed anymore.

Pseudocode

  # 1. Query PG table for bundle -> node mapping
  pg_table = ray.util.placement_group_table(pg)
  bundle_to_node = pg_table["bundles_to_node_id"]

  # 2. Get GPU bundles
  gpu_bundles = [(i, bundle_to_node[str(i)])
                 for i, b in enumerate(pg.bundle_specs)
                 if b.get("GPU", 0)]

  # 3. Sort: driver node first, then group by node
  driver_node = ray.get_runtime_context().get_node_id()
  gpu_bundles.sort(key=lambda x: (0 if x[1] == driver_node else 1, x[1]))

  # 4. Now assign rank = position in sorted list
  for rank, (bundle_id, node_id) in enumerate(gpu_bundles):
      actor = ray.remote(RayWorkerProc).options(
          num_gpus=1,
          scheduling_strategy=PlacementGroupSchedulingStrategy(
              placement_group=pg,
              placement_group_bundle_index=bundle_id,
          ),
          runtime_env={"env_vars": {
              # Tell Ray not to set CUDA_VISIBLE_DEVICES (same pattern as EEP)
              "RAY_EXPERIMENTAL_NOSET_CUDA_VISIBLE_DEVICES": "1",
          }},
      ).remote(rank=rank, local_rank=..., ...)

2.3 Worker interface

  @ray.remote
  class RayWorkerProc:                                                                                                              

      class ResponseStatus(Enum):
          SUCCESS = auto()
          FAILURE = auto()

      def __init__(
          self,
          vllm_config: VllmConfig,
          local_rank: int,
          rank: int,
          distributed_init_method: str,
          input_shm_handle: Handle,
          is_driver_worker: bool,
      ):
          ...

      def wait_for_init(self) -> dict:
          return {
              "status": "READY",
              "handle": self.worker_response_mq.export_handle(),
          }

      def run(self) -> None:
          try:
              # MQ barrier — identical to WorkerProc
              self.rpc_broadcast_mq.wait_until_ready()
              self.worker_response_mq.wait_until_ready()

              # Busy loop — identical to WorkerProc.worker_busy_loop
              self._worker_busy_loop()
          except Exception:
              logger.exception("RayWorkerProc failed.")
              raise
          finally:
              self.shutdown()

      def _worker_busy_loop(self) -> None:
          """Identical to WorkerProc.worker_busy_loop."""
          ...

      def _enqueue_output(self, output: Any) -> None:
          """Identical to WorkerProc.enqueue_output."""
          ...

      def shutdown(self) -> None:
          """Identical to WorkerProc.shutdown.
          No signal handler cleanup needed — Ray handles process teardown."""
          ...

Key differences from WorkerProc:

Aspect WorkerProc RayWorkerProc
Process Model Plain class instantiated inside mp.Process via static worker_main() Ray actor process managed by Ray
Readiness Signaling Uses mp.Pipe (ready_writer.send) Uses wait_for_init() + ray.get()
Execution Model Runs MQ barrier + busy loop inside worker_main() Splits logic into run() invoked via actor.run.remote()
Parent/Driver Death Handling Monitors parent death via death pipe EOF Relies on Ray GCS ownership -- actor is killed when driver dies
Shutdown Handling Uses SIGTERM / SIGINT signal handlers for graceful shutdown Relies on ray.kill()-- no signal handling required

2.4 Communication

2.4.1 Control plane – MessageQueue
  1. The driver creates broadcast MessageQueue.
  2. The broadcast MQ handle is delivered to actors via Ray constructor args.
  3. Each actor creates its own response MessageQueue and reports the handle back, and the driver connects a SUB socket back to each actor’s XPUB.
  4. Both sides call wait_until_ready. Writers (XPUB) block on zmq.recv() to track SUB subscription acknowledgments until all readers are registered. Readers (SUB) automatically emit subscription messages on connection. The broadcast MQ barrier is established first, followed by the response MQ barriers.
  5. collective_rpc flows through MessageQueue identically to MultiprocExecutor. No ray.remote calls on the hot path.
Image Image
2.4.2 Data plane – NCCL collective group
  1. Setup: During Worker.init_device. This is completely transparent to the executor, so no change is needed.
  2. Communication: Exactly the same as MultiprocExecutor.

2.5 Lifecycle and fault handling

Scenario MultiprocExecutor RayExecutorV2
Executor detects worker death Blocks on child process sentinel on a background thread. Once detected, marks executor as failed, shuts down, and notifies engine via failure_callback(). Uses ObjectRef as a sentinel. Child actor runs a busy loop blocking forever; caller uses ray.wait(). Once detected, follows same procedure as MultiprocExecutor. Existing pattern: busy loop + child death detection.
Worker detects executor death Child receives EOF error on parent crash through the death pipe. Implicit — if the owner actor dies, Ray GCS automatically kills all owned actors.
Graceful shutdown Closes death pipe, sends SIGTERM, then SIGKILL if necessary. Uses ray.kill(). Existing pattern: CoreEngineActorManager.close().

3. Test plan

Unit tests

  • Improve test coverage for MessageQueue cross-node TCP path
  • Add vllm/tests/distributed/test_ray_v2_executor.py mimicking vllm/tests/distributed/test_multiproc_executor.py
  • Adjust tests/distributed/test_multi_node_assignment.py::test_multi_node_assignment

Integration tests

  • Single-node tensor / pipeline parallelism
  • Cross-node tensor / pipeline parallelism
  • External placement group / bundle index compatibility

Benchmark
No overhead over MultiprocExecutor

4. Migration plan

Phase 1 (v0.x): Feature flag, disabled by default

Phase 2 (v0.x+2): Enabled by default
Set VLLM_USE_RAY_V2_EXECUTOR_BACKEND to True by default

Phase 3 (v0.x+4): Always enabled
Remove https://github.com/vllm-project/vllm/blob/main/vllm/v1/executor/ray_executor.py while maintaining backward compatibility.

Feedback Period

1 week.

CC List

@simon-mo @WoosukKwon @youkaichao @tlrmchlsmth @njhill @kouroshHakha

Any Other Things

No response

Before submitting a new issue...

  • Make sure you already searched for relevant issues, and asked the chatbot living at the bottom right corner of the documentation page, which can answer lots of frequently asked questions.

Activity

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Metadata

Metadata

Assignees

No one assigned

    Labels

    RFCstaleOver 90 days of inactivity

    Type

    No type

    Projects

    No projects

      Milestone

      No milestone

      Relationships

      None yet

      Development

      No branches or pull requests

      Issue actions