Ejemplos «Hello World» de Ray para el entorno de ejecución de IA

Importante

Esta característica está en versión preliminar pública.

Esta página tiene un ejemplo sencillo y funcional para cada una de las siguientes bibliotecas Ray en AI Runtime:

Prerequisites

La última CLI de Databricks instalada y autenticada. Consulte Instalación o actualización de la CLI de Databricks y Autenticación para la CLI de Databricks.

Arranque del clúster de Ray

Cuando envías una carga de trabajo, esta command se ejecuta en todos los nodos simultáneamente. Para usar Ray en varios nodos, el script de arranque usa NODE_RANK para decidir el rol de cada nodo. 0 inicia el nodo principal de Ray, y el resto se une como nodo de trabajo.

Cada ejemplo en esta página utiliza un compartido ray_bootstrap.sh para gestionar esta configuración. Solo necesitas una copia de este archivo junto con los scripts de ejemplo que estés ejecutando.

#!/bin/bash
# NODE_RANK=0 is the Ray head: it starts the cluster and runs the entrypoint
# script, then tears the cluster down. Every other rank joins as a worker and
# stays until the head goes away.
#
# The entrypoint to run on the head is passed via RAY_ENTRYPOINT, a path
# relative to CODE_SOURCE_PATH (e.g. "ray_train.py").
set -e

if [ -z "${RAY_ENTRYPOINT:-}" ]; then
    echo "RAY_ENTRYPOINT is not set; expected a script path relative to CODE_SOURCE_PATH." >&2
    exit 1
fi

RAY_HEAD_PORT=6379
GPUS_PER_NODE=${LOCAL_WORLD_SIZE:-1}

if [ "${NODE_RANK:-0}" = "0" ]; then
    echo "NODE_RANK=0: Starting Ray head node with $GPUS_PER_NODE GPU(s)..."
    ray start --head \
        --port=$RAY_HEAD_PORT \
        --num-gpus=$GPUS_PER_NODE \
        --dashboard-host=0.0.0.0

    # Always stop the cluster on exit, even if the entrypoint fails.
    trap 'ray stop' EXIT

    echo "Ray head node started. Running $RAY_ENTRYPOINT..."
    python "$CODE_SOURCE_PATH/$RAY_ENTRYPOINT"
else
    echo "NODE_RANK=$NODE_RANK: Connecting to Ray head at $MASTER_ADDR:$RAY_HEAD_PORT..."
    # Retry loop to wait for head to be ready. Note: omit --block, since it runs
    # forever and the head's `ray stop` only tears down local processes, leaving
    # the worker stuck. Without --block, `ray start` returns once this node joins
    # and we control our own exit below.
    joined=""
    for i in $(seq 1 12); do
        if ray start --address="$MASTER_ADDR:$RAY_HEAD_PORT" --num-gpus=$GPUS_PER_NODE 2>/dev/null; then
            joined=1
            break
        fi
        echo "Attempt $i failed, retrying in 5s..."
        sleep 5
    done
    if [ -z "$joined" ]; then
        echo "Worker failed to join the Ray head after all retries; aborting." >&2
        exit 1
    fi

    # `ray health-check` exits non-zero once the head runs `ray stop`, letting
    # this worker exit so the whole job can terminate. The counter backstops
    # against a hang.
    echo "Worker joined; waiting for the head to finish its work..."
    for _ in $(seq 1 360); do
        if ! ray health-check --address "$MASTER_ADDR:$RAY_HEAD_PORT" 2>/dev/null; then
            break
        fi
        sleep 5
    done
    echo "Head is no longer healthy; stopping local Ray and exiting."
    ray stop
fi

Cada ejemplo YAML invoca el bootstrap configurando RAY_ENTRYPOINT y llamando a ray_bootstrap.sh:

command: |
  cd $CODE_SOURCE_PATH
  RAY_ENTRYPOINT=ray_train.py bash ray_bootstrap.sh

LOCAL_WORLD_SIZE está configurado por AI Runtime al número de GPUs en cada nodo, por lo que GPUS_PER_NODE escala automáticamente con el tipo de GPU que solicites. AI Runtime establece MASTER_ADDR en la dirección IP del nodo principal, que los trabajadores utilizan para localizar y unirse al clúster de Ray.

Ray Core

El ejemplo muestra cómo programar el trabajo en cada GPU del clúster usando @ray.remote(num_gpus=1), que indica a Ray que coloque cada tarea en una GPU separada. Cada tarea indica en qué nodo y en qué GPU física se ejecutó, lo que confirma que las tareas se distribuyeron entre nodos en lugar de concentrarse en un único nodo.

Carga de trabajo YAML

ray_core.yaml solicita 2 nodos con 1 GPU A10 cada uno (GPU_1xA10), dando al clúster un total de 2 GPUs:

experiment_name: ray-core-example

environment:
  version: '5'
  dependencies:
    - ray[default]

code_source:
  type: snapshot
  snapshot:
    root_path: .

compute:
  num_accelerators: 2
  accelerator_type: GPU_1xA10

command: |
  cd $CODE_SOURCE_PATH
  RAY_ENTRYPOINT=ray_core.py bash ray_bootstrap.sh

max_retries: 0
timeout_minutes: 15
env_variables:
  NCCL_DEBUG: INFO

Script

ray_core.py despacha una tarea por GPU. Como Ray configura CUDA_VISIBLE_DEVICES en la única GPU asignada en cada tarea, current_device() siempre devuelve 0. El script utiliza ray.get_gpu_ids() y CUDA_VISIBLE_DEVICES para informar de la asignación física real:

@ray.remote(num_gpus=1)
def hello_from_gpu():
    node_rank = os.environ.get("NODE_RANK", "?")
    ray_gpu_ids = ray.get_gpu_ids()
    visible = os.environ.get("CUDA_VISIBLE_DEVICES", "")
    gpu_name = subprocess.run(
        ["nvidia-smi", "--query-gpu=name", "--format=csv,noheader"],
        capture_output=True, text=True, check=True,
    ).stdout.strip()
    return f"Hello from node {node_rank} | Ray GPU id {ray_gpu_ids} | CUDA_VISIBLE_DEVICES={visible} | {gpu_name}"

total_gpus = int(ray.cluster_resources().get("GPU", 0))
futures = [hello_from_gpu.remote() for _ in range(total_gpus)]
results = ray.get(futures)

El guion completo está en escrituras completas al final de esta página.

Envío de la ejecución

databricks air run -f ray_core.yaml --watch

Tren de Ray

El ejemplo entrena una pequeña MLP con datos sintéticos. prepare_model mueve el modelo a la GPU del trabajador y lo envuelve en DDP. prepare_data_loader añade un DistributedSampler para que cada trabajador vea un fragmento diferente de los datos, y ray.train.report devuelve métricas por época al controlador.

Carga de trabajo YAML

ray_train.yaml solicita 2 nodos con 1 GPU A10 cada uno. ray[train] Instala los extras de Ray Train:

experiment_name: ray-train-example

environment:
  version: '5'
  dependencies:
    - ray[train]
    - torch

code_source:
  type: snapshot
  snapshot:
    root_path: .

compute:
  num_accelerators: 2
  accelerator_type: GPU_1xA10

command: |
  cd $CODE_SOURCE_PATH
  RAY_ENTRYPOINT=ray_train.py bash ray_bootstrap.sh

max_retries: 0
timeout_minutes: 15
env_variables:
  NCCL_DEBUG: INFO

Escritura de entrenamiento

ray_train.py define un bucle de entrenamiento por trabajador y se TorchTrainer configura para usar todas las GPUs del clúster:

def train_loop_per_worker(config):
    model = nn.Sequential(nn.Linear(128, 256), nn.ReLU(), nn.Linear(256, 10))
    model = prepare_model(model)  # DDP wrap + move to this worker's GPU

    x = torch.randn(1024, 128)
    y = torch.randint(0, 10, (1024,))
    loader = DataLoader(TensorDataset(x, y), batch_size=64, shuffle=True)
    loader = prepare_data_loader(loader)  # adds DistributedSampler

    for epoch in range(config["epochs"]):
        ...
        ray.train.report({"epoch": epoch, "loss": epoch_loss / len(loader)})

trainer = TorchTrainer(
    train_loop_per_worker,
    train_loop_config={"lr": 1e-3, "epochs": 5},
    scaling_config=ScalingConfig(num_workers=total_gpus, use_gpu=True),
)
result = trainer.fit()

El guion completo está en escrituras completas al final de esta página.

Envío de la ejecución

databricks air run -f ray_train.yaml --watch

Datos de Ray

El ejemplo construye una canalización sintética: un map por fila añade características derivadas, un filter conserva solo las filas con número par y un map_batches aplica una transformación vectorizada de NumPy. Llamar a count() y sum() al final desencadena la ejecución.

Carga de trabajo YAML

ray_data.yaml Solicita 2 nodos. Los clústeres heterogéneos de CPU/GPU aún no están soportados en Ray Data on AI Runtime, por lo que este ejemplo mantiene la canalización en las CPUs. El GPU_1xA10 tipo de nodo determina el tamaño del clúster:

experiment_name: ray-data-example

environment:
  version: '5'
  dependencies:
    - ray[data]

code_source:
  type: snapshot
  snapshot:
    root_path: .

compute:
  num_accelerators: 2
  accelerator_type: GPU_1xA10

command: |
  cd $CODE_SOURCE_PATH
  RAY_ENTRYPOINT=ray_data.py bash ray_bootstrap.sh

max_retries: 0
timeout_minutes: 15

Procesamiento de scripts

ray_data.py define una tubería de tres etapas y imprime resultados agregados:

ds = ray.data.range(10_000)

def add_features(row):
    n = row["id"]
    return {"id": n, "squared": n * n, "is_even": n % 2 == 0}

def scale_batch(batch):
    batch["scaled"] = batch["squared"] * 0.001
    return batch

# Ray executes these stages in parallel across the cluster.
ds = ds.map(add_features)
ds = ds.filter(lambda row: row["is_even"])
ds = ds.map_batches(scale_batch, batch_format="numpy")

print(f"Pipeline produced {ds.count()} rows")
print(f"Sum of scaled feature: {ds.sum('scaled'):.2f}")

El guion completo está en escrituras completas al final de esta página.

Envío de la ejecución

databricks air run -f ray_data.yaml --watch

Ray Tune

El ejemplo ejecuta 8 pruebas, 4 a la vez en 4 GPUs. Cada prueba entrena una pequeña MLP con datos sintéticos usando una combinación seleccionada aleatoriamente de tasa de aprendizaje, tamaño de la capa oculta y tamaño del lote.

Carga de trabajo YAML

ray_tune.yaml solicita 4 nodos con 1 GPU A10 cada uno, dando 4 GPUs para hasta 4 pruebas concurrentes:

experiment_name: ray-tune-example

environment:
  version: '5'
  dependencies:
    - ray[tune]
    - torch

code_source:
  type: snapshot
  snapshot:
    root_path: .

compute:
  num_accelerators: 4
  accelerator_type: GPU_1xA10

command: |
  cd $CODE_SOURCE_PATH
  RAY_ENTRYPOINT=ray_tune.py bash ray_bootstrap.sh

max_retries: 0
timeout_minutes: 30

Ajuste del script

ray_tune.py configura el espacio de búsqueda y lanza 8 ensayos con ASHA, lo que detiene las pruebas que no rinden antes de tiempo:

tuner = tune.Tuner(
    tune.with_resources(train_fn, resources={"gpu": 1}),
    param_space={
        "lr": tune.loguniform(1e-4, 1e-1),
        "hidden_size": tune.choice([64, 128, 256]),
        "batch_size": tune.choice([32, 64, 128]),
    },
    tune_config=tune.TuneConfig(
        metric="loss",
        mode="min",
        scheduler=ASHAScheduler(max_t=20, grace_period=3, reduction_factor=2),
        num_samples=8,
    ),
)
results = tuner.fit()
best = results.get_best_result("loss", "min")
print(f"Best config: {best.config}")

tune.with_resources(train_fn, resources={"gpu": 1}) reserva una GPU por prueba. Con 4 GPUs, Ray Tune ejecuta 4 pruebas a la vez y comienza el siguiente lote cuando terminan las pruebas. El guion completo está en escrituras completas al final de esta página.

Envío de la ejecución

databricks air run -f ray_tune.yaml --watch

Inspeccionar una ejecución

Después de enviarla, puedes consultar los registros de estado y de streaming:

databricks air get <run-id>
databricks air logs <run-id>

databricks air logs transmite desde el nodo 0 por defecto, que es donde se ejecuta el controlador Ray. Para ver los registros de un nodo de trabajo, pasa --node 1, --node 2 y así sucesivamente.

Pasos siguientes

Guiones completos

ray_core.py

"""Ray Core remote-task example on AI Runtime.

Dispatches one @ray.remote task per GPU across the cluster. Each task prints
which node and physical GPU it was assigned to, confirming tasks reached every
node. Run after ray_bootstrap.sh has started the cluster.
"""

import os
import subprocess
import time

import ray

ray.init(address="auto")

num_nodes = int(os.environ.get("NUM_NODES", 1))
gpus_per_node = int(os.environ.get("LOCAL_WORLD_SIZE", 1))
expected_gpus = num_nodes * gpus_per_node

for _ in range(30):
    if len(ray.nodes()) >= num_nodes and ray.cluster_resources().get("GPU", 0) >= expected_gpus:
        break
    time.sleep(2)

total_gpus = int(ray.cluster_resources().get("GPU", 0))
if total_gpus < expected_gpus:
    raise SystemExit(
        f"Expected {expected_gpus} GPU(s) but Ray only sees {total_gpus}; " "check GPU discovery on all nodes."
    )

print(f"Ray cluster ready: {len(ray.nodes())} node(s), {total_gpus} GPU(s)")
print(f"Cluster resources: {ray.cluster_resources()}\n")


@ray.remote(num_gpus=1)
def hello_from_gpu():
    node_rank = os.environ.get("NODE_RANK", "?")
    # Ray sets CUDA_VISIBLE_DEVICES to the single assigned GPU, so
    # current_device() always returns 0. Report the physical GPU via
    # nvidia-smi and the Ray GPU ID instead.
    ray_gpu_ids = ray.get_gpu_ids()
    visible = os.environ.get("CUDA_VISIBLE_DEVICES", "")
    gpu_name = subprocess.run(
        ["nvidia-smi", "--query-gpu=name", "--format=csv,noheader"],
        capture_output=True,
        text=True,
        check=True,
    ).stdout.strip()
    return f"Hello from node {node_rank} | Ray GPU id {ray_gpu_ids} | CUDA_VISIBLE_DEVICES={visible} | {gpu_name}"


print(f"Launching {total_gpus} task(s), one per GPU across the cluster...")
futures = [hello_from_gpu.remote() for _ in range(total_gpus)]
results = ray.get(futures)

for r in results:
    print(r)

ray.shutdown()

ray_train.py

"""Ray Train distributed training example on AI Runtime.

Trains a small MLP on synthetic data with one training worker per GPU using
Ray Train's TorchTrainer. Ray Train places the workers across the cluster
(one per GPU) and wires up torch.distributed; the per-worker train loop just
uses `ray.train.torch` helpers to move the model/data to the right device.
"""

import os

import ray
import torch
import torch.nn as nn
from ray.train import ScalingConfig
from ray.train.torch import TorchTrainer, prepare_data_loader, prepare_model
from torch.utils.data import DataLoader, TensorDataset

# Connect to the cluster started by ray_bootstrap.sh.
ray.init(address="auto")

num_nodes = int(os.environ.get("NUM_NODES", 1))
total_gpus = int(ray.cluster_resources().get("GPU", 0))
if total_gpus < 1:
    raise SystemExit("No GPUs registered with Ray; check GPU discovery on the cluster.")
print(f"Cluster ready: {num_nodes} node(s), {total_gpus} GPU(s) available")
print(f"Launching a Ray Train run with {total_gpus} worker(s), one per GPU\n")


def train_loop_per_worker(config):
    """Runs on each Ray Train worker; one worker is pinned to one GPU."""
    # prepare_model wraps the model in DDP and moves it to this worker's GPU.
    model = nn.Sequential(nn.Linear(128, 256), nn.ReLU(), nn.Linear(256, 10))
    model = prepare_model(model)

    x = torch.randn(1024, 128)
    y = torch.randint(0, 10, (1024,))
    loader = DataLoader(TensorDataset(x, y), batch_size=64, shuffle=True)
    # prepare_data_loader shards the data across workers and moves batches to the GPU.
    loader = prepare_data_loader(loader)

    optimizer = torch.optim.Adam(model.parameters(), lr=config["lr"])
    loss_fn = nn.CrossEntropyLoss()

    for epoch in range(config["epochs"]):
        model.train()
        epoch_loss = 0.0
        for inputs, labels in loader:
            optimizer.zero_grad()
            loss = loss_fn(model(inputs), labels)
            loss.backward()
            optimizer.step()
            epoch_loss += loss.item()
        # ray.train.report surfaces metrics back to the driver.
        ray.train.report({"epoch": epoch, "loss": epoch_loss / len(loader)})


trainer = TorchTrainer(
    train_loop_per_worker,
    train_loop_config={"lr": 1e-3, "epochs": 5},
    scaling_config=ScalingConfig(num_workers=total_gpus, use_gpu=True),
)

result = trainer.fit()
# result.metrics holds the last reported dict (may be None if nothing was
# reported on the final iteration); fall back to a plain message.
print(f"\nTraining finished. Final metrics: {result.metrics or 'see per-worker logs above'}")

ray.shutdown()

ray_data.py

"""Ray Data distributed preprocessing example on AI Runtime.

Builds a Ray Dataset and runs a distributed map / map_batches / filter
pipeline across CPU actors spread over the cluster. On AI Runtime, Ray Data
runs on CPU actors (heterogeneous CPU/GPU clusters are not supported yet), so
this example deliberately keeps the transforms on CPU. The common shape is Ray
Data preprocessing feeding into a Ray Train run.
"""

import os

import ray

# Connect to the cluster started by ray_bootstrap.sh.
ray.init(address="auto")

num_nodes = int(os.environ.get("NUM_NODES", 1))
num_cpus = int(ray.cluster_resources().get("CPU", 0))
print(f"Cluster ready: {num_nodes} node(s), {num_cpus} CPU(s) available")

# A simple synthetic dataset; range() produces a distributed Ray Dataset.
ds = ray.data.range(10_000)


def add_features(row):
    """Per-row transform, runs distributed across CPU tasks."""
    n = row["id"]
    return {"id": n, "squared": n * n, "is_even": n % 2 == 0}


def scale_batch(batch):
    """Vectorized per-batch transform (numpy), more efficient than per-row."""
    batch["scaled"] = batch["squared"] * 0.001
    return batch


# Distributed pipeline: map -> filter -> map_batches, then aggregate.
ds = ds.map(add_features)
ds = ds.filter(lambda row: row["is_even"])
ds = ds.map_batches(scale_batch, batch_format="numpy")

count = ds.count()
total = ds.sum("scaled")
print(f"\nPipeline produced {count} rows (even numbers only)")
print(f"Sum of scaled feature: {total:.2f}")
print("\nSample of 5 processed rows:")
for row in ds.take(5):
    print(f"  {row}")

ray.shutdown()

ray_tune.py

"""Ray Tune hyperparameter search example on AI Runtime.

Runs 8 trials across all available GPUs in the cluster (one GPU per trial).
Uses ASHA scheduler to prune unpromising trials early.
"""

import os
import ray
import torch
import torch.nn as nn
from ray import tune
from ray.tune.schedulers import ASHAScheduler

ray.init(address="auto")

num_nodes = int(os.environ.get("NUM_NODES", 1))
total_gpus = int(ray.cluster_resources().get("GPU", 0))
if total_gpus < 1:
    raise SystemExit("No GPUs registered with Ray; check GPU discovery on the cluster.")
print(f"Cluster ready: {num_nodes} node(s), {total_gpus} GPU(s) available")
print(f"Running 8 trials with up to {total_gpus} in parallel\n")


def train_fn(config):
    """Single trial: trains a small MLP on synthetic data for one GPU."""
    device = torch.device("cuda" if torch.cuda.is_available() else "cpu")

    model = nn.Sequential(
        nn.Linear(128, config["hidden_size"]),
        nn.ReLU(),
        nn.Linear(config["hidden_size"], 10),
    ).to(device)

    optimizer = torch.optim.Adam(model.parameters(), lr=config["lr"])
    loss_fn = nn.CrossEntropyLoss()

    for epoch in range(20):
        x = torch.randn(config["batch_size"], 128, device=device)
        y = torch.randint(0, 10, (config["batch_size"],), device=device)

        optimizer.zero_grad()
        loss = loss_fn(model(x), y)
        loss.backward()
        optimizer.step()

        tune.report({"loss": loss.item(), "epoch": epoch})


tuner = tune.Tuner(
    tune.with_resources(train_fn, resources={"gpu": 1}),
    param_space={
        "lr": tune.loguniform(1e-4, 1e-1),
        "hidden_size": tune.choice([64, 128, 256]),
        "batch_size": tune.choice([32, 64, 128]),
    },
    tune_config=tune.TuneConfig(
        metric="loss",
        mode="min",
        scheduler=ASHAScheduler(max_t=20, grace_period=3, reduction_factor=2),
        num_samples=8,
    ),
)

results = tuner.fit()
best = results.get_best_result("loss", "min")
print(f"\nBest config: {best.config}")
print(f"Best loss:   {best.metrics['loss']:.4f}")

ray.shutdown()