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
- 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.
- Process cleanup: When using Ray as the external orchestrator with vLLM's mp backend, there will be zombie vLLM processes upon fault / OOM scenarios.
- 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.
- 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.
- Sort bundles: Driver node gets the lowest rank, and same node bundles get contiguous ranks.
- Compute global and per-node local rank.
- Classify local vs. remote for MessageQueue creation (more details in the following section).
- 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
- The driver creates broadcast
MessageQueue.
- The broadcast MQ handle is delivered to actors via Ray constructor args.
- 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.
- 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.
collective_rpc flows through MessageQueue identically to MultiprocExecutor. No ray.remote calls on the hot path.
2.4.2 Data plane – NCCL collective group
- Setup: During
Worker.init_device. This is completely transparent to the executor, so no change is needed.
- 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...
Summary
Replace the current
RayExecutordistributed executor backend which uses Compiled Graph with a newRayExecutorV2that launches workers as plain Ray remote actors and follows theMultiprocExecutorcommunication 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:
torch.distributedalready provides NCCL collectives natively. vLLM'sMultiprocExecutoruses 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.distributedAPIs 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
MultiprocExecutorgolden pathMultiprocExecutoris 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
MultiprocExecutorso 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 thatMultiprocExecutorcannot support. For example, SkyRL utilizesVLLM_BUNDLE_INDICESto schedule inference engine actor (vLLM) with a specific placement.Best reasons NOT to remove Ray from vLLM entirely
Proposed Change
1. Goals / Non-goals
1.1 Goal
MultiprocExecutorby sharing the same control and data plane communication1.2 Non-goals
2. Proposed changes
2.1 API
[Selected] Option 1: Inherit from
MultiprocExecutorImplement
RayExecutorV2as a child class ofMultiprocExecutorto reuse maximal code.MultiprocExecutor.MultiprocExecutorcould breakRayExecutorV2unexpectedly. We can mitigate this with comprehensive test coverage, and one of the goals of this RFC is to alignRayExecutorV2closely withMultiprocExecutor.Option 2: Independent RayExecutorV2
Implement an isolated
RayExecutorV2although there are substantial common pieces betweenRayExecutorV2andMultiprocExecutor, e.g.execute_model,collective_rpc, etc.MultiprocExecutor.Option 3: Extract common pieces to a new class
MessageQueueExecutorRefactor the duplicated logic in
MultiprocExecutorandRayExecutorV2into a unifiedMessageQueueExecutor.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.
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.
Pseudocode
2.3 Worker interface
Key differences from WorkerProc:
mp.Processvia staticworker_main()mp.Pipe(ready_writer.send)wait_for_init()+ray.get()worker_main()run()invoked viaactor.run.remote()SIGTERM/SIGINTsignal handlers for graceful shutdownray.kill()-- no signal handling required2.4 Communication
2.4.1 Control plane – MessageQueue
MessageQueue.MessageQueueand reports the handle back, and the driver connects a SUB socket back to each actor’s XPUB.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.collective_rpcflows throughMessageQueueidentically toMultiprocExecutor. Noray.remotecalls on the hot path.2.4.2 Data plane – NCCL collective group
Worker.init_device. This is completely transparent to the executor, so no change is needed.MultiprocExecutor.2.5 Lifecycle and fault handling
failure_callback().ObjectRefas a sentinel. Child actor runs a busy loop blocking forever; caller usesray.wait(). Once detected, follows same procedure asMultiprocExecutor. Existing pattern: busy loop + child death detection.SIGTERM, thenSIGKILLif necessary.ray.kill(). Existing pattern:CoreEngineActorManager.close().3. Test plan
Unit tests
vllm/tests/distributed/test_ray_v2_executor.pymimickingvllm/tests/distributed/test_multiproc_executor.pytests/distributed/test_multi_node_assignment.py::test_multi_node_assignmentIntegration tests
Benchmark
No overhead over
MultiprocExecutor4. Migration plan
Phase 1 (v0.x): Feature flag, disabled by default
VLLM_USE_RAY_V2_EXECUTOR_BACKEND(Falseby default) in https://github.com/vllm-project/vllm/blob/main/vllm/envs.pydistributed_executor_backed=”ray”and useRayExecutor(compiled graph backend) by defaultPhase 2 (v0.x+2): Enabled by default
Set
VLLM_USE_RAY_V2_EXECUTOR_BACKENDtoTrueby defaultPhase 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...