Aller au contenu principal

Exemples Ray hello world pour la CLI AI Runtime

info

Aperçu

Cette fonctionnalité est en Aperçu public.

Cette page contient un exemple fonctionnel simple pour chacune des bibliothèques Ray suivantes sur AI Runtime :

Prérequis

La CLI air installée et authentifiée. Voir l'installation de l'AI Runtime CLI.

Amorçage de cluster Ray

Lorsque vous soumettez une charge de travail, le command s'exécute simultanément sur chaque nœud. Pour utiliser Ray sur plusieurs nœuds, le script de bootstrap utilise NODE_RANK pour décider du rôle de chaque nœud. 0 lance le head Ray, et tout le reste rejoint en tant que Worker.

Chaque exemple sur cette page utilise un ray_bootstrap.sh partagé pour gérer cette configuration. Vous n'avez besoin que d'une seule copie de ce fichier avec les scripts d'exemple que vous exécutez.

Bash
#!/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

Chaque exemple YAML appelle l'amorçage en définissant RAY_ENTRYPOINT et en appelant ray_bootstrap.sh:

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

LOCAL_WORLD_SIZE est défini par AI Runtime sur le nombre de GPU sur chaque nœud, de sorte que GPUS_PER_NODE monte en charge automatiquement avec le type de GPU que vous demandez. MASTER_ADDR est défini par AI Runtime sur l'adresse IP du nœud principal, que les Worker utilisent pour localiser et rejoindre le cluster Ray.

Ray Core

L'exemple montre comment planifier le travail sur chaque GPU du cluster en utilisant @ray.remote(num_gpus=1), ce qui indique à Ray de placer chaque tâche sur un GPU distinct. Chaque tâche indique sur quel nœud et quel GPU physique elle a été exécutée, confirmant ainsi que les tâches ont été réparties sur plusieurs nœuds plutôt que regroupées sur un seul.

Workload YAML

ray_core.yaml demande 2 nœuds avec 1 GPU A10 chacun (GPU_1xA10), ce qui donne au cluster un total de 2 GPU :

YAML
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 répartit une tâche par GPU. Comme Ray définit CUDA_VISIBLE_DEVICES sur le GPU unique assigné à l'intérieur de chaque tâche, current_device() renvoie toujours 0. Le script utilise ray.get_gpu_ids() et CUDA_VISIBLE_DEVICES pour signaler l'assignation physique réelle :

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

Le script complet se trouve dans Scripts complets à la fin de cette page.

Soumettre l'exécution

Bash
air run -f ray_core.yaml --watch

Ray Train

L'exemple entraîne un petit MLP sur des données synthétiques. prepare_model déplace le modèle vers le GPU du worker et l'encapsule dans DDP. prepare_data_loader ajoute un DistributedSampler afin que chaque worker voie une partition différente des données, et ray.train.report renvoie les métriques par époque au driver.

Workload YAML

ray_train.yaml demande 2 nœuds avec 1 GPU A10 chacun. ray[train] installe les compléments Ray Train :

YAML
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

Script de formation

ray_train.py définit une boucle d'entraînement par worker et configure TorchTrainer pour utiliser tous les GPU du cluster :

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

Le script complet se trouve dans Scripts complets à la fin de cette page.

Soumettre l'exécution

Bash
air run -f ray_train.yaml --watch

Ray Data

L'exemple construit un pipeline synthétique : un map par ligne ajoute des caractéristiques dérivées, un filter ne conserve que les lignes paires, et un map_batches applique une transformation NumPy vectorisée. L'appel de count() et sum() à la fin Trigger l'exécution.

Workload YAML

ray_data.yaml demande 2 nœuds. Les clusters hétérogènes CPU/GPU ne sont pas encore pris en charge dans Ray Data sur AI Runtime, cet exemple conserve donc le pipeline sur les CPU. Le type de nœud GPU_1xA10 détermine la taille du cluster :

YAML
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

Script de traitement

ray_data.py définit un pipeline à trois étapes et affiche les résultats agrégés :

Python
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}")

Le script complet se trouve dans Scripts complets à la fin de cette page.

Soumettre l'exécution

Bash
air run -f ray_data.yaml --watch

Ray Tune

L'exemple exécute 8 essais, 4 à la fois sur 4 GPU. Chaque essai entraîne un petit MLP sur des données synthétiques avec une combinaison échantillonnée de taux d'apprentissage, de taille cachée et de taille de batch.

Workload YAML

ray_tune.yaml demande 4 nœuds avec 1 GPU A10 chacun, ce qui donne 4 GPU pour un maximum de 4 essais simultanés :

YAML
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

Script de réglage

ray_tune.py configure l'espace de recherche et lance 8 essais avec ASHA, qui arrête prématurément les essais peu performants :

Python
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}) réserve un GPU par test. Avec 4 GPU, Ray Tune exécute 4 tests à la fois et start le prochain batch à mesure que les tests se terminent. Le script complet se trouve dans Scripts complets à la fin de cette page.

Soumettre l'exécution

Bash
air run -f ray_tune.yaml --watch

Inspecter une exécution

Après la soumission, vous pouvez vérifier le statut et Stream les Logs :

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

air logs streams depuis le nœud 0 par default, là où le Driver Ray s'exécute. Pour afficher les Logs d'un nœud worker, transmettez --node 1, --node 2, et ainsi de suite.

Étapes suivantes

Scripts complets

ray_core.py

Python
"""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

Python
"""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={&quot;lr&quot;: 1e-3, &quot;epochs&quot;: 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

Python
"""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

Python
"""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={&quot;gpu&quot;: 1}),
param_space={
&quot;lr&quot;: tune.loguniform(1e-4, 1e-1),
&quot;hidden_size&quot;: tune.choice([64, 128, 256]),
&quot;batch_size&quot;: 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()