Skip to content

RayDistributedExecutor: usar Ray para llevar workers a múltiples nodos y tarjetas

源码版本v0.25.1

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 RayWorkerWrapper actor 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);luego collective_rpc("adjust_rank", args=(rerank_mapping,)) rectifica el rank.
  • Forward por Ray Compiled Graph: en lugar de llamar execute_model por RPC en cada step, se construye un forward_dag(ray_executor.py:527-620)que con PP=2/TP=4 encadena cuatro workers como InputNode -> TP group -> MultiOutputNode;los tensores intermedios viajan por with_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_VARS queda fuera para evitar contaminación entre workers.
  • execute_model partido en dos fases: con sampler activo, execute_model no corre, solo guarda self.scheduler_output, y se espera a sample_tokens para 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-97initialize_ray_cluster, obtiene el placement group, llama _init_workers_ray, e inicializa has_connector, uses_sampler y el placeholder scheduler_output.
  • shutdown:99-113forward_dag.teardown() y luego ray.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 usa num_gpus=num_gpus, otras plataformas usan resources={ray_device_key: num_gpus}.
  • rerank:217-258ray.get recoge 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 luego collective_rpc("adjust_rank", ...) rectifica rpc_rank.
  • init_worker + load_model:343-368 — tras preparar all_kwargs, se hace collective_rpc("init_worker") + collective_rpc("init_device") + collective_rpc("load_model") de una sola vez.
  • pp_tp_workers:370-378 — rellena self.pp_tp_workers con la malla PP×TP para luego construir el compiled DAG.
  • execute_model + sample_tokens:390-432 — dos fases: execute_model solo guarda scheduler_output, y el verdadero cómputo ocurre en sample_tokens vía _execute_dag.
  • _execute_dag:434-468 — en la primera invocación construye self.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 devuelve FutureWrapper.
  • collective_rpc:470-495 — serializa el method con cloudpickle.dumps, lo envía por execute_method.remote(...) a cada worker, y luego ray.get(ray_worker_outputs, timeout=timeout);con non_block=True devuelve FutureWrapper.
  • _compiled_ray_dag:527-620 — arma el DAG con InputNode + worker.execute_model_ray.bind(...); entre PP usa with_tensor_transport(transport) para pasar tensores; al final experimental_compile(...).
  • check_health:625-628 — en la versión actual hace return directo; queda un TODO para implementar health check real de actores.
  • RayWorkerWrapper:56-91 — hereda de WorkerWrapperBase y añade adjust_rank, get_node_ip, get_node_and_physical_gpu_ids, y execute_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:

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

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:

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

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_ray comprueba explícitamente n_nodes == n_ips, y si no, lanza RuntimeError sugiriendo revisar VLLM_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 que get_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_model ve self.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_installation exige 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, fuerza VLLM_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_health hace return directo(ray_executor.py:625-628);si un actor cae, hay que esperar a que ray.get relance la excepción, no hay heartbeat activo.
  • shutdown con kill duro: tras forward_dag.teardown(), se hace ray.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.