Ray hello worldのAIランタイムCLI例

Important

この機能は パブリック プレビュー段階です

このページには、AI Runtime上の以下のレイライブラリそれぞれの簡単な動作例があります。

前提条件

air CLI がインストールされ、認証されました。 AI ランタイム CLI のインストールを参照してください。

レイクラスターブートストラップ

ワークロードを提出すると、 command はすべてのノードで同時に動作します。 複数のノードでレイを使用する際、ブートストラップスクリプトは各ノードの役割を NODE_RANK で決定します。 0 レイの頭が起動し、他のすべてが働き手として加わります。

このページの各例は共有の ray_bootstrap.sh を使ってこのセットアップを処理しています。 このファイルは、実行している例スクリプトと一緒に1つのコピーだけで十分です。

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

各例のYAMLは、 RAY_ENTRYPOINT を設定して ray_bootstrap.shを呼び出してブートストラップを呼び出します。

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

LOCAL_WORLD_SIZE はAI Runtimeによって各ノードのGPU数に設定されており、 GPUS_PER_NODE はあなたが要求したGPUタイプに応じて自動的にスケーリングされます。 MASTER_ADDR はAI RuntimeによってヘッドノードのIPアドレスに設定されており、作業者はこれを使ってRayクラスタの位置を特定し参加します。

レイ・コア

この例は、クラスター内のすべてのGPUで作業をスケジュール化する方法を示しており、 @ray.remote(num_gpus=1)はRayに各タスクを別々のGPUに割り当てるよう指示しています。 各タスクはどのノードと物理GPUに到達したかを報告し、タスクがノード間に分散されていることを確認しました。

ワークロードYAM

ray_core.yaml 2つのノードにそれぞれ1つのA10 GPUを要求する(GPU_1xA10)ため、クラスタ全体は2つのGPUとなります:

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

スクリプト

ray_core.py GPUごとに1つのタスクを送ります。 Rayは各タスク内の単一の割り当てGPUに割り当てる CUDA_VISIBLE_DEVICES を設定するため、 current_device() は常に0を返します。 スクリプトは ray.get_gpu_ids()CUDA_VISIBLE_DEVICES を使って実際の物理的な割り当てを報告します。

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

完全な脚本はこのページの最後に「 Full scripts 」に掲載されています。

実行の送信

air run -f ray_core.yaml --watch

レイ・トレイン

この例は、合成データ上で小型MLPを訓練しています。 prepare_model モデルをワーカーのGPUに移動し、DDPでラップします。 prepare_data_loader 各ワーカーが異なるシャードのデータを見るように DistributedSampler を追加し、 ray.train.report はエポックごとの指標をドライバーに返します。

ワークロードYAM

ray_train.yaml 2つのノードにそれぞれ1つのA10 GPUを要求します。 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

トレーニングスクリプト

ray_train.py 作業者ごとのトレーニングループを定義し、 TorchTrainer をクラスタ内のすべてのGPUを使用するように設定します:

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

完全な脚本はこのページの最後に「 Full scripts 」に掲載されています。

実行の送信

air run -f ray_train.yaml --watch

レイ データ

この例は合成パイプラインを構築します。行 map ごとは導出特徴量を追加し、 filter は偶数行のみを保持し、 map_batches はベクトル化されたNumPy変換を適用します。 最後に count()sum() を呼び出すと実行がトリガーされます。

ワークロードYAM

ray_data.yaml リクエストは2ノードです。 AI Runtime上のRay Dataで異種CPU/GPUクラスターはまだサポートされていないため、この例ではパイプラインをCPUに留めています。 GPU_1xA10ノードの種類によってクラスタサイズが決まります:

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

処理スクリプト

ray_data.py 3段階のパイプラインを定義し、集計結果を印刷します:

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

完全な脚本はこのページの最後に「 Full scripts 」に掲載されています。

実行の送信

air run -f ray_data.yaml --watch

レイ・チューン

この例は4つのGPUで4つのトライアルを同時に実行します。 各試行は、学習率、隠れサイズ、バッチサイズのサンプリング組み合わせを用いた合成データ上で小規模なMLPを訓練します。

ワークロードYAM

ray_tune.yaml 4ノードにそれぞれ1つのA10 GPUを要求し、最大4つの同時試験で4つのGPUを提供します:

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

チューニングスクリプト

ray_tune.py ASHAで探索空間を設定し、8つの試験を開始。 ASHAは性能不振の試験を早期に停止させる:

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}) トライアルごとにGPUを1つ予約しています。 4つのGPUを使うと、Ray Tuneは一度に4つのトライアルを実行し、トライアルが終わると次のロットを始めます。 完全な脚本はこのページの最後に「 Full scripts 」に掲載されています。

実行の送信

air run -f ray_tune.yaml --watch

ランを点検する

提出後、ステータスや配信ログを確認できます:

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

air logs デフォルトでノード0からストリームを流し、そこでレイドライバーが動作します。 ワーカーノードからログを閲覧するには、 --node 1--node 2などを渡します。

次のステップ

全文脚本

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