RayDistributedExecutor: usar Ray para llevar workers a múltiples nodos y tarjetas
Responsabilidades
En escenarios monoplaca, UniProcExecutor basta con un proceso. Pero en cuanto TP>1 o se cruza entre nodos, hace falta comunicación entre procesos. RayDistributedExecutor(ray_executor.py:64-68) es la implementación con backend Ray: envuelve cada worker como un Ray actor, los reparte sobre las GPU indicadas por el placement group, y usa Ray Compiled Graph para construir un pipeline que encadena el forward. uses_ray = True y supports_pp = True(ray_executor.py:67-68)declaran que está orquestado por Ray y que soporta paralelismo de pipeline (pipeline parallelism).
Hereda de Executor igual que UniProcExecutor, sin tocar las firmas de execute_model/sample_tokens/profile/sleep de la base; solo reescribe _init_executor, collective_rpc, execute_model, check_health, shutdown(ray_executor.py:70-97). La diferencia está solo en cómo viaja la RPC: en mismo proceso es llamada directa, en Ray es worker.execute_method.remote(...)(ray_executor.py:484-489).
Motivación de diseño
- Un Ray actor por worker: una
RayWorkerWrapperactor por cada GPU(ray_executor.py:200-213);Ray se encarga de scheduling, tolerancia a fallos y contabilidad de recursos, y el ejecutor no hace spawn de procesos por su cuenta. - Placement group fija la topología: el actor se vincula a un bundle con
PlacementGroupSchedulingStrategy(ray_executor.py:192-196);la afinidad de GPU queda decidida antes del arranque y no deriva en runtime. - Reordenado priorizando el nodo local: el orden de creación de los actores lo decide Ray (puede ser aleatorio), por lo que tras crear todos se reordenan por IP, poniendo primero los del mismo nodo que el driver(
ray_executor.py:234-258);luegocollective_rpc("adjust_rank", args=(rerank_mapping,))rectifica el rank. - Forward por Ray Compiled Graph: en lugar de llamar
execute_modelpor RPC en cada step, se construye unforward_dag(ray_executor.py:527-620)que con PP=2/TP=4 encadena cuatro workers comoInputNode -> TP group -> MultiOutputNode;los tensores intermedios viajan porwith_tensor_transport(transport="nccl"|"shm")vía NCCL o memoria compartida(ray_executor.py:582-591). - Replicación automática de variables de entorno: las variables relacionadas con vLLM del proceso driver se empujan por lotes a cada worker con
update_environment_variables(ray_executor.py:305-327);WORKER_SPECIFIC_ENV_VARSqueda fuera para evitar contaminación entre workers. - execute_model partido en dos fases: con sampler activo,
execute_modelno corre, solo guardaself.scheduler_output, y se espera asample_tokenspara lanzar_execute_dag(ray_executor.py:390-432);así el cálculo del bitmask de grammar en el planificador y el sampler puede solaparse con el forward.
Archivos clave
RayDistributedExecutor._init_executor:64-97—initialize_ray_cluster, obtiene el placement group, llama_init_workers_ray, e inicializahas_connector,uses_samplery el placeholderscheduler_output.shutdown:99-113—forward_dag.teardown()y luegoray.kill(worker)por cada uno; el log indica que SIGTERM es esperado._init_workers_ray:143-213— por cada bundle_index,ray.remote(...)(RayWorkerWrapper).remote(rpc_rank=rank);en GPU usanum_gpus=num_gpus, otras plataformas usanresources={ray_device_key: num_gpus}.rerank:217-258—ray.getrecoge la IP de nodo de cada actor; se ordena por «mismo nodo que el driver primero → menos actores en el nodo primero → IP menor primero», y luegocollective_rpc("adjust_rank", ...)rectifica rpc_rank.init_worker + load_model:343-368— tras prepararall_kwargs, se hacecollective_rpc("init_worker")+collective_rpc("init_device")+collective_rpc("load_model")de una sola vez.pp_tp_workers:370-378— rellenaself.pp_tp_workerscon la malla PP×TP para luego construir el compiled DAG.execute_model + sample_tokens:390-432— dos fases:execute_modelsolo guardascheduler_output, y el verdadero cómputo ocurre ensample_tokensvía_execute_dag._execute_dag:434-468— en la primera invocación construyeself.forward_dag = self._compiled_ray_dag(...);forward_dag.execute((scheduler_output, grammar_output))devuelve refs; sin PP,refs[0].get()bloquea hasta el resultado; en modo no bloqueante devuelveFutureWrapper.collective_rpc:470-495— serializa el method concloudpickle.dumps, lo envía porexecute_method.remote(...)a cada worker, y luegoray.get(ray_worker_outputs, timeout=timeout);connon_block=TruedevuelveFutureWrapper._compiled_ray_dag:527-620— arma el DAG conInputNode+worker.execute_model_ray.bind(...); entre PP usawith_tensor_transport(transport)para pasar tensores; al finalexperimental_compile(...).check_health:625-628— en la versión actual hacereturndirecto; queda un TODO para implementar health check real de actores.RayWorkerWrapper:56-91— hereda deWorkerWrapperBasey añadeadjust_rank,get_node_ip,get_node_and_physical_gpu_ids, yexecute_method(envuelve excepciones con un log de «posible deadlock» antes de relanzar).
Flujo de datos
Durante el arranque, _init_executor pasa por _init_workers_ray, el tramo más largo: elige bundles, crea los actores RayWorkerWrapper según el placement group, y luego ray.get recoge la IP de nodo de cada actor para reordenar y reranquear:
# 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))En runtime, cuando EngineCore.step llega a execute_model, si uses_sampler=True y hay tokens en el lote, el ejecutor solo retiene scheduler_output y devuelve(ray_executor.py:401-407);es al invocarse sample_tokens cuando se llama de verdad a _execute_dag, que en la primera llamada construye el compiled DAG:
# 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])
...Las RPC del plano de control(initialize_from_config, compile_or_warm_up_model, determine_available_memory, etc.)van todas por collective_rpc, es decir, execute_method.remote(...) por cada worker(ray_executor.py:484-489),y luego ray.get para recoger. Esto y el plano de datos(forward por compiled DAG)son dos caminos independientes.
Límites y fallos
- La IP de nodo debe ser única:
_init_workers_raycomprueba explícitamenten_nodes == n_ips, y si no, lanzaRuntimeErrorsugiriendo revisarVLLM_HOST_IP(ray_executor.py:295-303). - Loopback en mono-nodo: con varios actores en el mismo nodo,
driver_ip = "127.0.0.1", evitando queget_ip()tome la tarjeta equivocada y rompa la comunicación entre workers(ray_executor.py:329-341). - Si se rompe el orden de la máquina de estados, error: si
execute_modelveself.scheduler_output is not None, lanza directamente «sample_tokens() must be called after execute_model() returns None»(ray_executor.py:395-399). - Versión y dependencias de Ray Compiled Graph:
_check_ray_cgraph_installationexige ray>=2.43.0,ray[cgraph]y cupy opcional(ray_executor.py:497-525);si no, hace fail-fast en el arranque. - TPU/XPU por memoria compartida:
_init_executor, al detectar TPU/XPU, fuerzaVLLM_USE_RAY_COMPILED_DAG_CHANNEL_TYPE=shm(ray_executor.py:73-75),para esquivar NCCL de NVIDIA. - check_health hoy es vacío:
RayDistributedExecutor.check_healthhacereturndirecto(ray_executor.py:625-628);si un actor cae, hay que esperar a queray.getrelance la excepción, no hay heartbeat activo. - shutdown con kill duro: tras
forward_dag.teardown(), se haceray.kill(worker)(ray_executor.py:111-112);el estado no persistido del worker se pierde, y la documentación advierte que el log de SIGTERM es esperado.
Resumen
RayDistributedExecutor es la implementación del ejecutor de v1 para escenarios multinode y multigpu. Envuelve cada worker como un Ray actor, fija la topología con un placement group, y construye un DAG de forward con Ray Compiled Graph que encadena TP/PP; la RPC del plano de control sigue yendo por execute_method del actor. La interfaz que expone a EngineCore es idéntica a la de UniProcExecutor, solo cambia el mecanismo de RPC por debajo. Cómo se abstrae el ejecutor, en /executor/executor;cómo corre el forward dentro del worker, en /worker/workery /worker/gpu-model-runner.
Véase la documentación oficial: documentación de vLLM · README.