Skip to content

在 AIStudio Ray 训练任务中使用 verl 运行双 Worker GRPO

注意

本文中的平台预置镜像地址是示例。使用前,请在镜像中心或当前实例的预置镜像列表中确认当前可用区提供了对应仓库及 tag;如果列表中不存在,请改用页面显示的可用镜像地址。替换 Registry 域名并不能保证相同仓库或 tag 在其他环境或可用区中存在。

本教程复用单 GPU 教程中的 DeepSeek-R1-Distill-Qwen-1.5B 和 GSM8K,通过 AIStudio 托管的 Ray 集群在两个 Worker 上运行 verl。四张 A100 共同完成两步全参数 GRPO,运行结果写入共享存储。

实战目标

完成本教程后,您可以检查一条小规模的多 Worker 训练链路:

  • AIStudio 创建 2 个 Ray 节点,每个节点提供 2 张 A100;
  • 两个节点从相同路径读取固定的 verl 代码、模型和 GSM8K 数据;
  • verl 创建一个 colocated ActorRollout group,并把 4 个 rank 分配到两个节点;
  • 两个训练 step 产生 rollout 和 TensorBoard 指标,FSDP2 保存 world_size=4 的 checkpoint;
  • 任务结束后,日志、拓扑记录和 checkpoint 仍可从共享存储读取。

这次运行用于检查 AIStudio 托管 Ray、跨节点共享存储和 verl 分布式执行。1.5B 模型不需要四张 GPU 才能训练,本教程也不比较不同拓扑的吞吐或成本。

场景信息

本教程使用以下配置:

  • verlv0.6.0@ddd86f527a4af75095e4677b02b5aa272913a088,从共享存储读取固定 commit;
  • 镜像:平台预置镜像 cr.infini-ai.com/infini-ai/verl:app-verl0.6-transformers4.56.1-sglang0.5.2-mcore0.13.0-te2.2
  • 模型:ModelScope 的 deepseek-ai/DeepSeek-R1-Distill-Qwen-1.5B@6fc93244f442ee2b5ab5c8000687ef5f7ffe1d03
  • 数据集:Hugging Face 的 openai/gsm8k@740312add88f781978c0658806c59bc2815b9866,转换为 verl v0.6 使用的 Parquet 格式;
  • 资源来源:Spot;每次创建或重跑使用新的 RUN_ID
  • 分布式框架:Ray,由 AIStudio 创建和管理 Ray Head 与 Worker;
  • Worker 拓扑:2 个 Worker、每个 Worker 2 张 A100-SXM4-80GB,共 4 张 GPU;
  • 网络:RDMA 关闭,Ray Worker 使用 NCCL_IB_DISABLE=1
  • 训练范围:全参数 GRPO,FSDP2,train_batch_size=4,每个 prompt 生成 2 个回答,共运行 2 个 step;
  • 共享存储:训练任务中的挂载路径为 /mnt/verl-reproduction,输出写入 runs/verl/two-worker-grpo/<RUN_ID>/

本场景在 AIStudio 中怎样运行

AIStudio 的分布式框架决定 Ray 集群由谁创建:

  • 单机:适合一个 Worker。verl 可以在 Worker 内启动本地 Ray,但不会使用第二个节点。
  • Ray(本教程):适合多个 Worker。AIStudio 先创建 Ray Head 和 Worker,再在 Head 中执行 Driver 启动命令;启动脚本只连接现有集群,不执行 ray startray job submit

两个 AIStudio Worker 同时也是两个 Ray 节点。verl 根据 trainer.nnodes=2trainer.n_gpus_per_node=2 创建一个 world size 为 4 的 colocated ActorRollout group:

language-text
AIStudio Ray 训练任务:2 个 Worker × 2 张 GPU

Worker 0 / Ray Head
├── verl Driver
└── 2 个 ActorRollout actor / 2 张 GPU

Worker 1 / Ray Worker
└── 2 个 ActorRollout actor / 2 张 GPU

同置(colocated)的 ActorRollout group 在同一组 rank 中完成模型加载、SGLang rollout 和 FSDP2 actor 更新。与把 rollout、reward 和 actor 拆成独立阶段后分别绑定到节点的流水线不同,本教程检查的是一个跨两个节点的 worker group。训练完成后,校验脚本从 Gloo 和 NCCL 日志确认 rank 03 组成同一个通信组、两个节点各运行 2 个 WorkerDict 进程,并检查 checkpoint 的完整分片。rank 编号由运行时分配。

两个 Worker 从 /mnt/verl-reproduction 读取同一份代码、模型和数据,并把输出写入同一共享目录。AIStudio 托管 Ray 的 Driver 写法和资源规则见提交 Ray Driver 入口命令

开始前准备

在创建 GPU 任务前,使用挂载同一共享存储的 AICoder 或开发机检查固定输入,并保存本教程的校验脚本和启动脚本。本文在准备环境中使用 /mnt/rlinf-reproduction;如果您的 AICoder 使用其他挂载路径,只需替换本节的 PREP_ROOT。训练任务中的所有路径仍统一使用 /mnt/verl-reproduction

复用单 GPU 教程的固定输入

本教程复用单 GPU verl GRPO 教程准备的代码、模型和数据。共享存储中应已经包含:

language-text
<shared-volume-root>/
├── code/verl/ddd86f527a4af75095e4677b02b5aa272913a088/
├── models/modelscope/deepseek-ai/DeepSeek-R1-Distill-Qwen-1.5B/
│   └── 6fc93244f442ee2b5ab5c8000687ef5f7ffe1d03/
└── datasets/huggingface/openai/gsm8k/
    └── 740312add88f781978c0658806c59bc2815b9866/
        └── processed/verl-v0.6.0-v1/

如果这些目录尚未准备,请先完成单 GPU 教程中的固定 verl 代码模型下载和校验以及 GSM8K 转换。双 Worker 任务直接读取同一份文件,不需要复制模型权重。

选用平台预置镜像

创建训练任务时选择以下平台预置镜像,无需导入到自己的镜像仓库:

language-text
cr.infini-ai.com/infini-ai/verl:app-verl0.6-transformers4.56.1-sglang0.5.2-mcore0.13.0-te2.2

该镜像提供 PyTorch、Ray、Transformers 和 SGLang 等运行依赖。verl 源代码从共享存储中的固定 checkout 加载,任务运行时不联网克隆代码或安装 Python 依赖。

保存多 Worker 校验脚本

把以下内容保存为 /mnt/rlinf-reproduction/tools/verl/two-worker-grpo/verify-two-worker-grpo-v1.py。该脚本在训练前检查两个节点、GPU、共享存储和固定输入,并在训练后根据持久日志检查 ActorRollout rank 的节点分布、Socket 通信、指标、rollout 和 4 个 rank 的 checkpoint 分片。

verify-two-worker-grpo-v1.pyPython28 KiB下载原始文件
显示代码隐藏代码,文件 verify-two-worker-grpo-v1.py653 行
python
#!/usr/bin/env python3
"""Positive preflight and completion checks for the bounded two-Worker verl GRPO run."""

from __future__ import annotations

import argparse
import hashlib
import importlib
import importlib.metadata
import json
import math
import os
import re
import subprocess
import sys
from datetime import datetime, timezone
from pathlib import Path
from typing import Any


EXPECTED_COMMIT = "ddd86f527a4af75095e4677b02b5aa272913a088"
EXPECTED_DATASET_REVISION = "740312add88f781978c0658806c59bc2815b9866"
EXPECTED_DATA_CHECKSUMS = "0150c256df1500d5edf5f5462434b287ada29a9674ddf1d59ef779629d6e8a4e"
EXPECTED_MODEL_MANIFEST = "0ea1d330342b4b9efbd1c3648360fbc4c2e3b1d5abd6120e49ef837649001f7b"
EXPECTED_GPU_NAME = "NVIDIA A100-SXM4-80GB"
MIN_CUDA_12_DRIVER = (525, 60, 13)
REQUIRED_SCALARS = (
    "training/global_step",
    "actor/pg_loss",
    "actor/grad_norm",
    "actor/lr",
    "critic/score/mean",
)
FATAL_PATTERNS = (
    r"Traceback \(most recent call last\)",
    r"Error executing job with overrides",
    r"CUDA out of memory",
    r"OutOfMemoryError",
    r"RayTaskError",
    r"ActorDiedError",
    r"WorkerCrashedError",
    r"\bNCCL (?:WARN|ERROR)\b",
    r"\bnccl(?:UnhandledCudaError|SystemError|InternalError|InvalidArgument|InvalidUsage|RemoteError)\b",
    r"\bProcessGroupNCCL\b[^\n]*(?:abort|watchdog|timeout|timed out)",
    r"Segmentation fault",
    r"Killed(?:\s|$)",
)


def sha256(path: Path) -> str:
    digest = hashlib.sha256()
    with path.open("rb") as stream:
        for chunk in iter(lambda: stream.read(1024 * 1024), b""):
            digest.update(chunk)
    return digest.hexdigest()


def run(*command: str) -> str:
    return subprocess.run(command, check=True, text=True, capture_output=True).stdout.strip()


def package_record(distribution: str, module_name: str | None = None) -> dict[str, object]:
    module = importlib.import_module(module_name or distribution.replace("-", "_"))
    try:
        distribution_version = importlib.metadata.version(distribution)
    except importlib.metadata.PackageNotFoundError:
        distribution_version = None
    return {
        "distribution_version": distribution_version,
        "module_version": getattr(module, "__version__", None),
        "module_file": getattr(module, "__file__", None),
    }


def atomic_json(path: Path, data: dict[str, object]) -> None:
    temporary = path.with_suffix(path.suffix + ".tmp")
    temporary.write_text(json.dumps(data, indent=2, ensure_ascii=False) + "\n", encoding="utf-8")
    os.replace(temporary, path)


def driver_version_tuple(version: str) -> tuple[int, int, int]:
    parts = version.split(".")
    if not parts or any(not part.isdigit() for part in parts):
        raise RuntimeError(f"Unrecognized NVIDIA driver version: {version}")
    numbers = [int(part) for part in parts[:3]]
    numbers.extend([0] * (3 - len(numbers)))
    return tuple(numbers)  # type: ignore[return-value]


def ordered_nodes(ray: Any) -> list[dict[str, Any]]:
    nodes = [node for node in ray.nodes() if node.get("Alive")]
    current_node_id = str(ray.get_runtime_context().get_node_id())
    head = [node for node in nodes if str(node.get("NodeID")) == current_node_id]
    if len(head) != 1:
        raise RuntimeError(f"Cannot identify the Ray Head: current={current_node_id} nodes={nodes}")
    workers = sorted(
        (node for node in nodes if node is not head[0]),
        key=lambda node: str(node.get("NodeManagerAddress", "")),
    )
    return head + workers


def preflight(args: argparse.Namespace) -> int:
    import ray
    from ray.util.scheduling_strategies import NodeAffinitySchedulingStrategy

    output_dir = args.output_dir.resolve()
    output_dir.mkdir(parents=True, exist_ok=True)
    sentinel = output_dir / "shared-storage-sentinel"
    sentinel.write_text(f"{args.run_id}\n", encoding="utf-8")
    ray.init(address="auto", ignore_reinit_error=True)
    nodes = ordered_nodes(ray)
    if len(nodes) != args.expected_nodes:
        raise RuntimeError(f"Expected {args.expected_nodes} live Ray nodes, found {len(nodes)}")

    cluster_resources = ray.cluster_resources()
    expected_total_gpus = args.expected_nodes * args.expected_gpus_per_node
    if int(cluster_resources.get("GPU", 0)) != expected_total_gpus:
        raise RuntimeError(
            f"Expected {expected_total_gpus} cluster GPUs, found {cluster_resources}"
        )
    for index, node in enumerate(nodes):
        node_gpus = int(node.get("Resources", {}).get("GPU", 0))
        if node_gpus != args.expected_gpus_per_node:
            raise RuntimeError(f"Node {index} exposes {node_gpus} GPUs: {node}")

    @ray.remote(num_cpus=1)
    def probe_node(
        ordinal: int,
        storage_root: str,
        verl_root: str,
        model_root: str,
        data_root: str,
        output_root: str,
        expected_node_id: str,
    ) -> dict[str, object]:
        import json as remote_json
        import socket
        from pathlib import Path as RemotePath

        import pyarrow.parquet as pq
        import ray as remote_ray

        def remote_run(*command: str) -> str:
            return subprocess.run(command, check=True, text=True, capture_output=True).stdout.strip()

        def read_one(path: RemotePath) -> int:
            with path.open("rb") as stream:
                return len(stream.read(1))

        storage = RemotePath(storage_root)
        source = RemotePath(verl_root)
        model = RemotePath(model_root)
        data = RemotePath(data_root)
        output = RemotePath(output_root)
        actual_node_id = str(remote_ray.get_runtime_context().get_node_id())
        if actual_node_id != expected_node_id:
            raise RuntimeError(
                f"Node {ordinal}: affinity target {expected_node_id}, ran on {actual_node_id}"
            )
        sentinel_path = output / "shared-storage-sentinel"
        if sentinel_path.read_text(encoding="utf-8") != f"{args.run_id}\n":
            raise RuntimeError(f"Node {ordinal}: cannot read the shared-storage sentinel")
        marker_root = output / "preflight-nodes"
        marker_root.mkdir(parents=True, exist_ok=True)
        marker_path = marker_root / f"node-{ordinal}.json"
        marker_path.write_text(
            remote_json.dumps({"ordinal": ordinal, "node_id": actual_node_id}) + "\n",
            encoding="utf-8",
        )
        source_commit = remote_run("git", "-C", str(source), "rev-parse", "HEAD")
        if source_commit != EXPECTED_COMMIT:
            raise RuntimeError(f"Node {ordinal}: unexpected verl commit {source_commit}")
        if remote_run("git", "-C", str(source), "status", "--porcelain"):
            raise RuntimeError(f"Node {ordinal}: the verl checkout is not clean")

        model_manifest = sha256(model / "SHA256SUMS")
        if model_manifest != EXPECTED_MODEL_MANIFEST:
            raise RuntimeError(f"Node {ordinal}: unexpected model manifest {model_manifest}")
        dataset_manifest_path = data / "dataset-manifest.json"
        dataset_manifest = remote_json.loads(dataset_manifest_path.read_text(encoding="utf-8"))
        if dataset_manifest.get("revision") != EXPECTED_DATASET_REVISION:
            raise RuntimeError(f"Node {ordinal}: unexpected dataset revision")
        data_checksums = sha256(data / "SHA256SUMS")
        if data_checksums != EXPECTED_DATA_CHECKSUMS:
            raise RuntimeError(f"Node {ordinal}: unexpected dataset checksums")
        datasets: dict[str, dict[str, object]] = {}
        for split, expected_rows in (("train", 7473), ("test", 1319)):
            path = data / f"{split}.parquet"
            rows = pq.read_metadata(path).num_rows
            if rows != expected_rows or read_one(path) != 1:
                raise RuntimeError(f"Node {ordinal}: invalid {split} parquet")
            datasets[split] = {"path": str(path), "rows": rows, "bytes": path.stat().st_size}
        if read_one(model / "model.safetensors") != 1:
            raise RuntimeError(f"Node {ordinal}: model weights are not readable")

        gpu_rows = remote_run(
            "nvidia-smi",
            "--query-gpu=driver_version,name,memory.total,uuid",
            "--format=csv,noheader,nounits",
        ).splitlines()
        parsed_gpus = []
        for row in gpu_rows:
            driver, name, memory, uuid = (part.strip() for part in row.split(",", 3))
            parsed_gpus.append(
                {"driver": driver, "name": name, "memory_mib": int(memory), "uuid": uuid}
            )
        if len(parsed_gpus) != args.expected_gpus_per_node:
            raise RuntimeError(f"Node {ordinal}: unexpected GPUs {parsed_gpus}")
        if any(gpu["name"] != EXPECTED_GPU_NAME for gpu in parsed_gpus):
            raise RuntimeError(f"Node {ordinal}: unexpected GPU model {parsed_gpus}")

        mount = remote_json.loads(
            remote_run("findmnt", "-J", "-T", str(storage), "-o", "TARGET,SOURCE,FSTYPE,OPTIONS")
        )["filesystems"][0]
        if mount["fstype"] == "overlay" or "rw" not in mount["options"].split(","):
            raise RuntimeError(f"Node {ordinal}: storage is not a writable shared mount {mount}")
        if os.environ.get("NCCL_IB_DISABLE") != "1":
            raise RuntimeError(f"Node {ordinal}: NCCL_IB_DISABLE must be 1")

        packages = {
            "torch": package_record("torch"),
            "ray": package_record("ray"),
            "transformers": package_record("transformers"),
            "sglang": package_record("sglang"),
            "pyarrow": package_record("pyarrow"),
            "tensorboard": package_record("tensorboard"),
            "verl": package_record("verl"),
        }
        return {
            "ordinal": ordinal,
            "node_id": actual_node_id,
            "node_ip": remote_ray.util.get_node_ip_address(),
            "hostname": socket.gethostname(),
            "gpus": parsed_gpus,
            "mount": mount,
            "python": {"version": sys.version, "executable": sys.executable},
            "packages": packages,
            "verl_source": {"path": str(source), "commit": source_commit, "clean": True},
            "model": {"path": str(model), "manifest_sha256": model_manifest},
            "dataset": {
                "path": str(data),
                "checksums_sha256": data_checksums,
                "manifest": dataset_manifest,
                "splits": datasets,
            },
        }

    probes = []
    for ordinal, node in enumerate(nodes):
        strategy = NodeAffinitySchedulingStrategy(node_id=str(node["NodeID"]), soft=False)
        probes.append(
            probe_node.options(scheduling_strategy=strategy).remote(
                ordinal,
                str(args.storage_root),
                str(args.verl_root),
                str(args.model_root),
                str(args.data_root),
                str(output_dir),
                str(node["NodeID"]),
            )
        )
    workers = ray.get(probes)
    drivers = {gpu["driver"] for worker in workers for gpu in worker["gpus"]}
    model_manifests = {worker["model"]["manifest_sha256"] for worker in workers}
    dataset_manifests = {worker["dataset"]["checksums_sha256"] for worker in workers}
    package_identities = {
        json.dumps(
            {"python": worker["python"], "packages": worker["packages"]},
            sort_keys=True,
            default=str,
        )
        for worker in workers
    }
    if min(driver_version_tuple(version) for version in drivers) < MIN_CUDA_12_DRIVER:
        raise RuntimeError(f"Unsupported NVIDIA driver versions: {sorted(drivers)}")
    if len(model_manifests) != 1 or len(dataset_manifests) != 1:
        raise RuntimeError("Input identity differs across Ray nodes")
    if len(package_identities) != 1:
        raise RuntimeError("Python package identities differ across Ray nodes")
    marker_root = output_dir / "preflight-nodes"
    marker_records = [
        json.loads((marker_root / f"node-{ordinal}.json").read_text(encoding="utf-8"))
        for ordinal in range(args.expected_nodes)
    ]
    if {record["node_id"] for record in marker_records} != {
        worker["node_id"] for worker in workers
    }:
        raise RuntimeError("The Driver cannot read both per-node shared-storage markers")

    for worker in workers:
        print(
            "WORKER_TOPOLOGY "
            + json.dumps(
                {
                    "ordinal": worker["ordinal"],
                    "node_id": worker["node_id"],
                    "node_ip": worker["node_ip"],
                    "hostname": worker["hostname"],
                    "gpu_names": [gpu["name"] for gpu in worker["gpus"]],
                    "driver": worker["gpus"][0]["driver"],
                    "mount_source": worker["mount"]["source"],
                },
                sort_keys=True,
            )
        )

    record = {
        "status": "passed",
        "phase": "preflight",
        "checked_at": datetime.now(timezone.utc).isoformat(),
        "run_id": args.run_id,
        "target_steps": args.target_steps,
        "train_batch_size": args.train_batch_size,
        "rollout_n": args.rollout_n,
        "expected_nodes": args.expected_nodes,
        "expected_gpus_per_node": args.expected_gpus_per_node,
        "cluster_resources": cluster_resources,
        "driver_versions": sorted(drivers),
        "ray_context": {
            "driver_node_id": str(ray.get_runtime_context().get_node_id()),
            "address": os.environ.get("RAY_ADDRESS", "auto"),
        },
        "workers": workers,
        "shared_storage_markers": marker_records,
    }
    atomic_json(output_dir / "preflight.json", record)
    print(json.dumps(record, indent=2, ensure_ascii=False))
    ray.shutdown()
    return 0


def load_scalars(tensorboard_root: Path) -> tuple[list[Path], dict[str, list[dict[str, float | int]]]]:
    from tensorboard.backend.event_processing.event_accumulator import EventAccumulator

    event_files = sorted(tensorboard_root.rglob("events.out.tfevents.*"))
    if not event_files or any(path.stat().st_size == 0 for path in event_files):
        raise RuntimeError("No non-empty TensorBoard event file was produced")
    scalars: dict[str, list[dict[str, float | int]]] = {}
    for event_file in event_files:
        accumulator = EventAccumulator(str(event_file), size_guidance={"scalars": 0})
        accumulator.Reload()
        for tag in accumulator.Tags().get("scalars", []):
            scalars.setdefault(tag, []).extend(
                {"step": event.step, "value": event.value, "wall_time": event.wall_time}
                for event in accumulator.Scalars(tag)
            )
    return event_files, scalars


def actor_placement_from_log(
    log_text: str,
    preflight_record: dict[str, object],
    world_size: int,
    run_id: str,
) -> dict[str, object]:
    workers = preflight_record.get("workers", [])
    if not isinstance(workers, list):
        raise RuntimeError("Preflight worker records are invalid")
    worker_by_hostname = {
        str(worker["hostname"]): worker
        for worker in workers
        if isinstance(worker, dict) and worker.get("hostname")
    }
    plain_log = re.sub(r"\x1b\[[0-9;]*m", "", log_text)
    nccl_pattern = re.compile(
        r"\(WorkerDict pid=(?P<prefix_pid>\d+)(?:, ip=[^)]+)?\)\s+"
        r"(?P<hostname>[^:\s]+):(?P<pid>\d+):\d+\s+\[\d+\] NCCL INFO "
        r"ncclCommInitRankConfig comm \S+ rank (?P<rank>\d+) nranks (?P<world>\d+) "
        r"cudaDev \d+ nvmlDev (?P<nvml_dev>\d+) busId (?P<bus_id>\S+) "
        r"commId (?P<comm_id>\S+) - Init (?P<phase>START|COMPLETE)"
    )
    communicators: dict[str, dict[int, dict[str, dict[str, object]]]] = {}
    for match in nccl_pattern.finditer(plain_log):
        if int(match.group("world")) != world_size:
            continue
        prefix_pid = int(match.group("prefix_pid"))
        pid = int(match.group("pid"))
        if prefix_pid != pid:
            raise RuntimeError(f"WorkerDict/NCCL PID mismatch: {prefix_pid} != {pid}")
        event: dict[str, object] = {
            "rank": int(match.group("rank")),
            "hostname": match.group("hostname"),
            "pid": pid,
            "nvml_dev": int(match.group("nvml_dev")),
            "bus_id": match.group("bus_id"),
        }
        comm_id = match.group("comm_id")
        rank = int(event["rank"])
        phase = match.group("phase")
        phases = communicators.setdefault(comm_id, {}).setdefault(rank, {})
        existing = phases.get(phase)
        if existing is not None and existing != event:
            raise RuntimeError(f"Conflicting NCCL {phase} evidence for rank {rank}")
        phases[phase] = event

    candidates: list[tuple[str, list[dict[str, object]]]] = []
    for comm_id, ranks in communicators.items():
        if sorted(ranks) != list(range(world_size)):
            continue
        actors = []
        valid = True
        for rank in range(world_size):
            phases = ranks[rank]
            if set(phases) != {"START", "COMPLETE"} or phases["START"] != phases["COMPLETE"]:
                valid = False
                break
            actors.append(dict(phases["START"]))
        if valid and len({(actor["hostname"], actor["pid"]) for actor in actors}) == world_size:
            candidates.append((comm_id, actors))
    if not candidates:
        raise RuntimeError("No complete world-size ActorRollout NCCL communicator was found")
    baseline = [
        (actor["rank"], actor["hostname"], actor["pid"], actor["nvml_dev"], actor["bus_id"])
        for actor in candidates[0][1]
    ]
    if any(
        [
            (actor["rank"], actor["hostname"], actor["pid"], actor["nvml_dev"], actor["bus_id"])
            for actor in actors
        ]
        != baseline
        for _, actors in candidates[1:]
    ):
        raise RuntimeError("Multiple full-world NCCL communicators have conflicting rank placement")
    comm_ids = sorted(comm_id for comm_id, _ in candidates)
    actors = candidates[0][1]

    gloo_pattern = re.compile(
        r"\(WorkerDict pid=(?P<pid>\d+)(?:, ip=[^)]+)?\).*?\[Gloo\] "
        r"Rank (?P<rank>\d+) is connected to (?P<peers>\d+) peer ranks\. "
        r"Expected number of connected peer ranks is : (?P<expected>\d+)"
    )
    gloo_evidence: dict[tuple[int, int], set[tuple[int, int]]] = {}
    for match in gloo_pattern.finditer(plain_log):
        key = (int(match.group("pid")), int(match.group("rank")))
        gloo_evidence.setdefault(key, set()).add(
            (int(match.group("peers")), int(match.group("expected")))
        )

    node_counts: dict[str, int] = {}
    for actor in actors:
        rank = int(actor["rank"])
        pid = int(actor["pid"])
        hostname = str(actor["hostname"])
        worker = worker_by_hostname.get(hostname)
        if worker is None:
            raise RuntimeError(f"ActorRollout rank {rank} ran on unknown host {hostname}")
        peer_evidence = gloo_evidence.get((pid, rank), set())
        expected_peers = (world_size - 1, world_size - 1)
        if expected_peers not in peer_evidence:
            raise RuntimeError(
                f"Missing full-mesh Gloo evidence for rank {rank}: {sorted(peer_evidence)}"
            )
        node_ip = str(worker["node_ip"])
        prefix = rf"\(WorkerDict pid={pid}(?:, ip=[^)]+)?\)\s+{re.escape(hostname)}:{pid}:\d+"
        required_transport = (
            rf"{prefix}.*?NCCL INFO NCCL_IB_DISABLE set by environment to 1\.",
            rf"{prefix}.*?NCCL INFO NET/Socket : Using .*?{re.escape(node_ip)}<\d+>",
            rf"{prefix}.*?NCCL INFO Using network Socket",
        )
        if any(not re.search(pattern, plain_log) for pattern in required_transport):
            raise RuntimeError(f"Missing socket-only NCCL evidence for rank {rank}")
        actor.update(
            {
                "class_name": "WorkerDict",
                "node_id": str(worker["node_id"]),
                "node_ip": node_ip,
                "connected_peer_ranks": expected_peers[0],
                "nccl_init": "complete",
                "nccl_transport": "Socket",
                "rdma_disabled": True,
            }
        )
        node_id = str(actor["node_id"])
        node_counts[node_id] = node_counts.get(node_id, 0) + 1
    expected_node_ids = {
        str(worker["node_id"])
        for worker in workers
        if isinstance(worker, dict) and worker.get("node_id")
    }
    if set(node_counts) != expected_node_ids or sorted(node_counts.values()) != [2, 2]:
        raise RuntimeError(f"Unexpected ActorRollout node distribution: {node_counts}")

    return {
        "status": "passed",
        "phase": "distributed_rank_placement",
        "source": "training_log",
        "checked_at": datetime.now(timezone.utc).isoformat(),
        "run_id": run_id,
        "expected_actors": world_size,
        "comm_ids": comm_ids,
        "node_counts": node_counts,
        "actors": actors,
    }


def postflight(args: argparse.Namespace) -> int:
    output_dir = args.output_dir.resolve()
    log_text = args.training_log.read_text(encoding="utf-8", errors="replace")
    for pattern in FATAL_PATTERNS:
        if re.search(pattern, log_text, flags=re.IGNORECASE):
            raise RuntimeError(f"Fatal training marker found: {pattern}")
    required_log_markers = (
        "[validate_config] All configuration checks passed successfully!",
        f"Total training steps: {args.target_steps}",
        f"global_step_{args.target_steps}",
    )
    missing_markers = [marker for marker in required_log_markers if marker not in log_text]
    if missing_markers:
        raise RuntimeError(f"Training log is missing completion markers: {missing_markers}")

    preflight_record = json.loads((output_dir / "preflight.json").read_text(encoding="utf-8"))
    workers = preflight_record.get("workers", [])
    if len(workers) != args.expected_nodes:
        raise RuntimeError("Preflight topology evidence is missing")
    actor_placement = actor_placement_from_log(
        log_text, preflight_record, args.world_size, args.run_id
    )
    atomic_json(output_dir / "actor-placement.json", actor_placement)

    resolved_command = (output_dir / "resolved-command.txt").read_text(encoding="utf-8")
    required_overrides = (
        f"trainer.n_gpus_per_node={args.expected_gpus_per_node}",
        f"trainer.nnodes={args.expected_nodes}",
        "+ray_kwargs.ray_init.address=auto",
        "+ray_kwargs.ray_init.runtime_env.env_vars.TENSORBOARD_DIR=",
        "actor_rollout_ref.rollout.name=sglang",
        "actor_rollout_ref.actor.strategy=fsdp2",
    )
    missing_overrides = [item for item in required_overrides if item not in resolved_command]
    if missing_overrides:
        raise RuntimeError(f"Resolved command is missing topology/runtime overrides: {missing_overrides}")

    event_files, scalars = load_scalars(output_dir / "tensorboard")
    metric_summary = {}
    for tag in REQUIRED_SCALARS:
        events = scalars.get(tag, [])
        if not events:
            raise RuntimeError(f"TensorBoard scalar is missing: {tag}")
        latest = max(events, key=lambda event: int(event["step"]))
        if int(latest["step"]) < args.target_steps or not math.isfinite(float(latest["value"])):
            raise RuntimeError(f"TensorBoard scalar did not reach a finite target step: {tag}={latest}")
        metric_summary[tag] = latest

    expected_rollouts = args.train_batch_size * args.rollout_n
    rollout_summary = {}
    for step in range(1, args.target_steps + 1):
        path = output_dir / "rollouts" / f"{step}.jsonl"
        rows = [json.loads(line) for line in path.read_text(encoding="utf-8").splitlines() if line.strip()]
        if len(rows) != expected_rollouts:
            raise RuntimeError(f"Unexpected rollout count at step {step}: {len(rows)}")
        for row in rows:
            if row.get("step") != step or not row.get("input") or not isinstance(row.get("output"), str):
                raise RuntimeError(f"Invalid rollout row at step {step}")
            if not math.isfinite(float(row["score"])):
                raise RuntimeError(f"Non-finite rollout score at step {step}")
        rollout_summary[str(step)] = {
            "path": str(path),
            "rows": len(rows),
            "bytes": path.stat().st_size,
            "score_min": min(float(row["score"]) for row in rows),
            "score_max": max(float(row["score"]) for row in rows),
        }

    checkpoint_root = output_dir / "checkpoints"
    checkpoint = checkpoint_root / f"global_step_{args.target_steps}"
    actor = checkpoint / "actor"
    fsdp_config_path = actor / "fsdp_config.json"
    fsdp_config = json.loads(fsdp_config_path.read_text(encoding="utf-8"))
    if fsdp_config != {"FSDP_version": 2, "world_size": args.world_size}:
        raise RuntimeError(f"Unexpected FSDP checkpoint configuration: {fsdp_config}")
    checkpoint_files = {}
    for kind in ("model", "optim", "extra_state"):
        for rank in range(args.world_size):
            path = actor / f"{kind}_world_size_{args.world_size}_rank_{rank}.pt"
            if not path.is_file() or path.stat().st_size == 0:
                raise RuntimeError(f"Checkpoint shard is missing or empty: {path}")
            checkpoint_files[str(path.relative_to(output_dir))] = path.stat().st_size
    for path in (
        checkpoint / "data.pt",
        checkpoint_root / "latest_checkpointed_iteration.txt",
        fsdp_config_path,
    ):
        if not path.is_file() or path.stat().st_size == 0:
            raise RuntimeError(f"Checkpoint file is missing or empty: {path}")
        checkpoint_files[str(path.relative_to(output_dir))] = path.stat().st_size
    if (checkpoint_root / "latest_checkpointed_iteration.txt").read_text().strip() != str(args.target_steps):
        raise RuntimeError("Checkpoint tracker does not match the target step")

    record = {
        "status": "passed",
        "phase": "postflight",
        "verified_at": datetime.now(timezone.utc).isoformat(),
        "run_id": args.run_id,
        "target_steps": args.target_steps,
        "train_batch_size": args.train_batch_size,
        "rollout_n": args.rollout_n,
        "world_size": args.world_size,
        "topology": {
            "nodes": len(workers),
            "gpus_per_node": preflight_record["expected_gpus_per_node"],
            "workers": workers,
            "actor_placement": actor_placement,
        },
        "tensorboard": {
            "event_files": [{"path": str(path), "bytes": path.stat().st_size} for path in event_files],
            "metrics": metric_summary,
        },
        "rollouts": rollout_summary,
        "checkpoint": {
            "path": str(checkpoint),
            "fsdp_config": fsdp_config,
            "files": checkpoint_files,
        },
    }
    atomic_json(output_dir / "verification.json", record)
    (output_dir / "SUCCESS").write_text(datetime.now(timezone.utc).isoformat() + "\n", encoding="utf-8")
    print(json.dumps(record, indent=2, ensure_ascii=False))
    return 0


def parse_args() -> argparse.Namespace:
    parser = argparse.ArgumentParser()
    subparsers = parser.add_subparsers(dest="phase", required=True)
    common = argparse.ArgumentParser(add_help=False)
    common.add_argument("--run-id", required=True)
    common.add_argument("--output-dir", type=Path, required=True)
    common.add_argument("--target-steps", type=int, required=True)
    common.add_argument("--train-batch-size", type=int, required=True)
    common.add_argument("--rollout-n", type=int, required=True)

    before = subparsers.add_parser("preflight", parents=[common])
    before.add_argument("--storage-root", type=Path, required=True)
    before.add_argument("--verl-root", type=Path, required=True)
    before.add_argument("--model-root", type=Path, required=True)
    before.add_argument("--data-root", type=Path, required=True)
    before.add_argument("--expected-nodes", type=int, required=True)
    before.add_argument("--expected-gpus-per-node", type=int, required=True)

    after = subparsers.add_parser("postflight", parents=[common])
    after.add_argument("--training-log", type=Path, required=True)
    after.add_argument("--expected-nodes", type=int, required=True)
    after.add_argument("--expected-gpus-per-node", type=int, required=True)
    after.add_argument("--world-size", type=int, required=True)
    return parser.parse_args()


if __name__ == "__main__":
    parsed = parse_args()
    handlers = {"preflight": preflight, "postflight": postflight}
    selected_handler = handlers[parsed.phase]
    raise SystemExit(selected_handler(parsed))

保存多 Worker 训练启动脚本

把以下内容保存为 /mnt/rlinf-reproduction/tools/verl/two-worker-grpo/run-two-worker-grpo-v1.sh。脚本连接 AIStudio 已创建的 Ray 集群,使用 trainer.nnodes=2trainer.n_gpus_per_node=2 运行两步 GRPO,并在完成检查通过后退出。

run-two-worker-grpo-v1.shBash6.6 KiB下载原始文件
显示代码隐藏代码,文件 run-two-worker-grpo-v1.sh183 行
bash
#!/usr/bin/env bash
set -euo pipefail

: "${STORAGE_ROOT:?Set STORAGE_ROOT to the shared-storage mount path}"
: "${RUN_ID:?Set a unique RUN_ID}"

readonly VERL_COMMIT=ddd86f527a4af75095e4677b02b5aa272913a088
readonly DATASET_REVISION=740312add88f781978c0658806c59bc2815b9866
readonly TARGET_STEPS=2
readonly TRAIN_BATCH_SIZE=4
readonly ROLLOUT_N=2
readonly EXPECTED_NODES=2
readonly GPUS_PER_NODE=2
readonly WORLD_SIZE=$((EXPECTED_NODES * GPUS_PER_NODE))
readonly VERL_ROOT="${STORAGE_ROOT}/code/verl/${VERL_COMMIT}"
readonly MODEL_ROOT="${STORAGE_ROOT}/models/modelscope/deepseek-ai/DeepSeek-R1-Distill-Qwen-1.5B/6fc93244f442ee2b5ab5c8000687ef5f7ffe1d03"
readonly DATA_ROOT="${STORAGE_ROOT}/datasets/huggingface/openai/gsm8k/${DATASET_REVISION}/processed/verl-v0.6.0-v1"
readonly TOOL_ROOT="${STORAGE_ROOT}/tools/verl/two-worker-grpo"
readonly VERIFY_SCRIPT="${TOOL_ROOT}/verify-two-worker-grpo-v1.py"
readonly EXPECTED_VERIFY_SHA256="3a2d6ffc85beef7479b08762f4dc00820957ae21dc153227b981bf5eb1b7a527"
readonly OUTPUT_ROOT="${STORAGE_ROOT}/runs/verl/two-worker-grpo"
readonly OUTPUT_DIR="${OUTPUT_ROOT}/${RUN_ID}"

[[ "${RUN_ID}" =~ ^[A-Za-z0-9][A-Za-z0-9._-]*$ ]] || {
  echo "RUN_ID contains unsupported characters" >&2
  exit 1
}
if [[ -e "${OUTPUT_DIR}" ]]; then
  echo "Refusing to reuse OUTPUT_DIR: ${OUTPUT_DIR}" >&2
  exit 1
fi
mkdir -p "${OUTPUT_ROOT}"
mkdir "${OUTPUT_DIR}"

readonly DRIVER_LOG="${OUTPUT_DIR}/driver.log"
readonly TRAINING_LOG="${OUTPUT_DIR}/training.log"
readonly CHECKPOINT_ROOT="${OUTPUT_DIR}/checkpoints"
readonly TENSORBOARD_DIR="${OUTPUT_DIR}/tensorboard"
readonly ROLLOUT_DIR="${OUTPUT_DIR}/rollouts"
readonly HYDRA_DIR="${OUTPUT_DIR}/hydra"
exec > >(tee -a "${DRIVER_LOG}") 2>&1

on_exit() {
  local status=$?
  if (( status != 0 )) && [[ ! -e "${OUTPUT_DIR}/SUCCESS" ]]; then
    printf '%s exit=%s\n' "$(date -u +%Y-%m-%dT%H:%M:%SZ)" "${status}" > "${OUTPUT_DIR}/FAILED"
  fi
  exit "${status}"
}
trap on_exit EXIT

for path in "${VERL_ROOT}" "${MODEL_ROOT}" "${DATA_ROOT}"; do
  test -d "${path}"
done
for path in "${DATA_ROOT}/train.parquet" "${DATA_ROOT}/test.parquet" "${VERIFY_SCRIPT}"; do
  test -r "${path}"
done
echo "${EXPECTED_VERIFY_SHA256}  ${VERIFY_SCRIPT}" | sha256sum -c -

if [[ -n "${PYTHONPATH:-}" ]]; then
  export PYTHONPATH="${VERL_ROOT}:${PYTHONPATH}"
else
  export PYTHONPATH="${VERL_ROOT}"
fi
export PYTHONUNBUFFERED=1
export HYDRA_FULL_ERROR=1
export HF_HUB_OFFLINE=1
export HF_DATASETS_OFFLINE=1
export TRANSFORMERS_OFFLINE=1
export WANDB_MODE=disabled
export TOKENIZERS_PARALLELISM=true
export NCCL_IB_DISABLE=1
export CUDA_DEVICE_MAX_CONNECTIONS=1
export TENSORBOARD_DIR

echo "RUN_ID: ${RUN_ID}"
echo "STORAGE_ROOT: ${STORAGE_ROOT}"
echo "RAY_ADDRESS: ${RAY_ADDRESS:-auto}"
findmnt -T "${STORAGE_ROOT}" -o TARGET,SOURCE,FSTYPE,OPTIONS
nvidia-smi --query-gpu=driver_version,name,memory.total --format=csv,noheader
echo "VERL_COMMIT: $(git -C "${VERL_ROOT}" rev-parse HEAD)"
test -z "$(git -C "${VERL_ROOT}" status --porcelain)"
(cd "${MODEL_ROOT}" && sha256sum -c SHA256SUMS)
(cd "${DATA_ROOT}" && sha256sum -c SHA256SUMS)

python3 "${VERIFY_SCRIPT}" preflight \
  --run-id "${RUN_ID}" \
  --output-dir "${OUTPUT_DIR}" \
  --target-steps "${TARGET_STEPS}" \
  --train-batch-size "${TRAIN_BATCH_SIZE}" \
  --rollout-n "${ROLLOUT_N}" \
  --storage-root "${STORAGE_ROOT}" \
  --verl-root "${VERL_ROOT}" \
  --model-root "${MODEL_ROOT}" \
  --data-root "${DATA_ROOT}" \
  --expected-nodes "${EXPECTED_NODES}" \
  --expected-gpus-per-node "${GPUS_PER_NODE}"

training_command=(
  python3 -m verl.trainer.main_ppo
  "algorithm.adv_estimator=grpo"
  "algorithm.use_kl_in_reward=false"
  "data.train_files=${DATA_ROOT}/train.parquet"
  "data.val_files=${DATA_ROOT}/test.parquet"
  "data.train_batch_size=${TRAIN_BATCH_SIZE}"
  "data.max_prompt_length=512"
  "data.max_response_length=256"
  "data.filter_overlong_prompts=true"
  "data.truncation=error"
  "data.shuffle=false"
  "actor_rollout_ref.model.path=${MODEL_ROOT}"
  "actor_rollout_ref.model.use_shm=false"
  "actor_rollout_ref.model.use_remove_padding=true"
  "actor_rollout_ref.model.enable_gradient_checkpointing=true"
  "actor_rollout_ref.actor.strategy=fsdp2"
  "actor_rollout_ref.ref.strategy=fsdp2"
  "critic.strategy=fsdp2"
  "reward_model.strategy=fsdp2"
  "actor_rollout_ref.actor.optim.lr=1e-6"
  "actor_rollout_ref.actor.ppo_mini_batch_size=${TRAIN_BATCH_SIZE}"
  "actor_rollout_ref.actor.ppo_micro_batch_size_per_gpu=2"
  "actor_rollout_ref.actor.use_kl_loss=false"
  "actor_rollout_ref.actor.entropy_coeff=0"
  "actor_rollout_ref.actor.fsdp_config.param_offload=true"
  "actor_rollout_ref.actor.fsdp_config.optimizer_offload=true"
  "actor_rollout_ref.rollout.name=sglang"
  "actor_rollout_ref.rollout.tensor_model_parallel_size=1"
  "actor_rollout_ref.rollout.gpu_memory_utilization=0.3"
  "actor_rollout_ref.rollout.n=${ROLLOUT_N}"
  "actor_rollout_ref.rollout.log_prob_micro_batch_size_per_gpu=2"
  "actor_rollout_ref.rollout.max_num_seqs=16"
  "actor_rollout_ref.rollout.max_model_len=768"
  "actor_rollout_ref.rollout.max_num_batched_tokens=2048"
  "actor_rollout_ref.rollout.enable_chunked_prefill=false"
  "actor_rollout_ref.rollout.enforce_eager=true"
  "trainer.logger=[console,tensorboard]"
  "trainer.project_name=verl_aistudio"
  "trainer.experiment_name=${RUN_ID}"
  "trainer.n_gpus_per_node=${GPUS_PER_NODE}"
  "trainer.nnodes=${EXPECTED_NODES}"
  "trainer.val_before_train=false"
  "trainer.test_freq=-1"
  "trainer.save_freq=${TARGET_STEPS}"
  "trainer.total_epochs=1"
  "trainer.total_training_steps=${TARGET_STEPS}"
  "trainer.resume_mode=disable"
  "trainer.default_local_dir=${CHECKPOINT_ROOT}"
  "trainer.rollout_data_dir=${ROLLOUT_DIR}"
  "+ray_kwargs.ray_init.address=auto"
  "+ray_kwargs.ray_init.runtime_env.env_vars.TENSORBOARD_DIR=${TENSORBOARD_DIR}"
  "hydra.job.chdir=false"
  "hydra.output_subdir=null"
  "hydra.run.dir=${HYDRA_DIR}"
)

printf '%q ' "${training_command[@]}" > "${OUTPUT_DIR}/resolved-command.txt"
printf '\n' >> "${OUTPUT_DIR}/resolved-command.txt"
printf 'Executing: '
printf '%q ' "${training_command[@]}"
printf '\n'

set +e
"${training_command[@]}" 2>&1 | tee -a "${TRAINING_LOG}"
training_status=${PIPESTATUS[0]}
set -e
if (( training_status != 0 )); then
  echo "verl training exited with status ${training_status}" >&2
  exit "${training_status}"
fi

python3 "${VERIFY_SCRIPT}" postflight \
  --run-id "${RUN_ID}" \
  --output-dir "${OUTPUT_DIR}" \
  --target-steps "${TARGET_STEPS}" \
  --train-batch-size "${TRAIN_BATCH_SIZE}" \
  --rollout-n "${ROLLOUT_N}" \
  --training-log "${TRAINING_LOG}" \
  --expected-nodes "${EXPECTED_NODES}" \
  --expected-gpus-per-node "${GPUS_PER_NODE}" \
  --world-size "${WORLD_SIZE}"

trap - EXIT
echo "Verification passed: ${OUTPUT_DIR}/verification.json"

运行这个场景

按照以下步骤检查输入、创建 Ray 训练任务并确认分布式输出。

Step 1 校验共享存储中的固定输入

分配 GPU 前,在挂载同一共享存储的 AICoder 或开发机中执行:

language-bash
set -euo pipefail

export PREP_ROOT=/mnt/rlinf-reproduction
export VERL_COMMIT=ddd86f527a4af75095e4677b02b5aa272913a088
export MODEL_REVISION=6fc93244f442ee2b5ab5c8000687ef5f7ffe1d03
export DATASET_REVISION=740312add88f781978c0658806c59bc2815b9866
export VERL_ROOT="$PREP_ROOT/code/verl/$VERL_COMMIT"
export MODEL_ROOT="$PREP_ROOT/models/modelscope/deepseek-ai/DeepSeek-R1-Distill-Qwen-1.5B/$MODEL_REVISION"
export DATA_ROOT="$PREP_ROOT/datasets/huggingface/openai/gsm8k/$DATASET_REVISION/processed/verl-v0.6.0-v1"

findmnt -T "$PREP_ROOT" -o TARGET,SOURCE,FSTYPE,OPTIONS
test -w "$PREP_ROOT"

test "$(git -C "$VERL_ROOT" rev-parse HEAD)" = "$VERL_COMMIT"
test -z "$(git -C "$VERL_ROOT" status --porcelain)"
(cd "$MODEL_ROOT" && sha256sum -c SHA256SUMS)
(cd "$DATA_ROOT" && sha256sum -c SHA256SUMS)

命令应显示可写共享存储挂载;代码工作树保持干净,模型和数据的 SHA256SUMS 全部通过。

Step 2 校验多 Worker 脚本

在同一个准备环境中执行:

language-bash
set -euo pipefail

export TOOL_ROOT=/mnt/rlinf-reproduction/tools/verl/two-worker-grpo

python3 -m py_compile "$TOOL_ROOT/verify-two-worker-grpo-v1.py"
bash -n "$TOOL_ROOT/run-two-worker-grpo-v1.sh"
test "$(sha256sum "$TOOL_ROOT/verify-two-worker-grpo-v1.py" | cut -d ' ' -f 1)" = \
  3a2d6ffc85beef7479b08762f4dc00820957ae21dc153227b981bf5eb1b7a527
test "$(sha256sum "$TOOL_ROOT/run-two-worker-grpo-v1.sh" | cut -d ' ' -f 1)" = \
  70926158cbe15773023512f9369a1a94abc57abfbeedf5e24f482d0fa4a13a3a

两个文件都通过后,训练任务会执行与本文一致的版本化脚本。

Step 3 创建双 Worker Ray 训练任务

在 AIStudio 中选择 训练任务,然后设置:

  1. 资源类型:选择 spot
  2. 可用区:本文使用宁夏 B;也可以选择满足镜像、GPU 和共享存储条件的其他可用区。
  3. 镜像:选择本教程指定的 verl 0.6 + SGLang 平台预置镜像。
  4. Worker 规格:选择每个 Worker 包含 2 张 A100-SXM4-80GB 的规格。
  5. Worker 数量:填写 2
  6. 分布式框架:选择 Ray
  7. RDMA 配置:保持关闭。

这组资源用于检查两个 Ray 节点和 4 个 rank 的分布式输出。改变 Worker 数量或每个 Worker 的 GPU 数量时,需要同时修改训练参数和完成检查。

Step 4 配置 RuntimeEnv

RuntimeEnv 配置栏中填写:

language-json
{
  "working_dir": "/mnt/verl-reproduction/code/verl/ddd86f527a4af75095e4677b02b5aa272913a088",
  "excludes": [
    "/.git/",
    "/docs/",
    "/tests/"
  ],
  "env_vars": {
    "PYTHONPATH": "/mnt/verl-reproduction/code/verl/ddd86f527a4af75095e4677b02b5aa272913a088",
    "HF_HUB_OFFLINE": "1",
    "HF_DATASETS_OFFLINE": "1",
    "TRANSFORMERS_OFFLINE": "1",
    "WANDB_MODE": "disabled",
    "TOKENIZERS_PARALLELISM": "true",
    "NCCL_IB_DISABLE": "1",
    "NCCL_DEBUG": "INFO",
    "NCCL_CUMEM_ENABLE": "0",
    "CUDA_DEVICE_MAX_CONNECTIONS": "1",
    "RAY_DEDUP_LOGS": "0",
    "PYTHONUNBUFFERED": "1"
  }
}

working_dirPYTHONPATH 让 Ray 任务和 actor 加载同一份固定版本的 verl 源码;离线变量阻止训练期间联网下载;NCCL_IB_DISABLE=1 与未启用 RDMA 的任务配置保持一致。不要在 RuntimeEnv 中设置 RAY_ADDRESS、Head 地址、端口或任务 ID,这些值由 AIStudio 管理。

Step 5 配置共享存储和 TensorBoard

在存储配置中选择准备模型和数据时使用的同一共享存储卷,并设置:

  • 将该共享存储卷在 Worker 内的访问路径设置为 /mnt/verl-reproduction
  • 访问权限:可读写;
  • 共享存储至少预留 25 GiB,用于本次运行的日志和 4 个 rank 的 checkpoint 分片;
  • 环境变量:添加一个新的运行标识,例如 RUN_ID=verl-two-worker-grpo-exp-001

开启 任务可视化,在 日志存储路径 中填写:

language-text
/mnt/verl-reproduction/runs/verl/two-worker-grpo/${RUN_ID}/tensorboard

/mnt/rlinf-reproduction/mnt/verl-reproduction 可以是同一存储卷在准备环境和训练任务中的不同访问路径。目录中的相对路径保持不变,因此两个 Worker 读取的是同一份代码、模型和数据。

每次创建、克隆或重跑任务时更换 RUN_ID${RUN_ID} 的变量替换、固定绝对路径和任务结束后的查看方式见使用训练任务托管的 TensorBoard 服务

Step 6 填写 Driver 启动命令

启动命令 中填写:

language-bash
set -euo pipefail
export STORAGE_ROOT=/mnt/verl-reproduction
nvidia-smi --query-gpu=driver_version,name --format=csv,noheader
sha256sum /mnt/verl-reproduction/tools/verl/two-worker-grpo/run-two-worker-grpo-v1.sh | grep -q ^70926158cbe15773023512f9369a1a94abc57abfbeedf5e24f482d0fa4a13a3a
exec /bin/bash /mnt/verl-reproduction/tools/verl/two-worker-grpo/run-two-worker-grpo-v1.sh

启动命令只打印 NVIDIA 驱动版本和 GPU 型号、校验版本化启动脚本,然后执行共享存储中的 Shell 文件。Ray 集群和 RuntimeEnv 已由 AIStudio 准备,启动命令不创建 Ray 集群,也不临时生成 Python 文件。

Step 7 核对配置并创建任务

提交前检查:

  • Spot、2 个 Worker、每个 Worker 2 张 A100-SXM4-80GB、Ray、RDMA 关闭;
  • 平台预置的 verl 0.6 + SGLang 镜像;
  • 两个 Worker 的共享存储路径均为 /mnt/verl-reproduction,并且可读写;
  • RuntimeEnv、TensorBoard 和启动命令都使用 /mnt/verl-reproduction
  • RUN_ID 在本次任务中唯一;
  • 启动命令执行 run-two-worker-grpo-v1.sh,脚本哈希与 Step 2 一致。

确认后单击 确认创建。任务创建完成后,在任务详情中再次核对 Worker 数量、每个 Worker 的 GPU 数量、共享存储路径和 RUN_ID

Step 8 确认两个 Ray 节点和 ActorRollout 放置

打开任务详情的 任务日志。训练开始前,日志应显示:

  1. cluster_resources 中有 4 张 GPU,preflight.json 中有 2 个 Worker 记录。
  2. 两条 WORKER_TOPOLOGY 分别对应 Ray Head 和另一个 Ray Worker。
  3. 每条记录包含 2 张 NVIDIA A100-SXM4-80GB,并打印该节点的 NVIDIA 驱动版本;校验要求每个节点不低于 525.60.13,节点间版本不必完全相同。
  4. 两个节点都通过共享存储读写检查,并加载相同版本的 verl 源码、模型和数据清单。

训练完成后,校验脚本会从持久日志中找到一个完成初始化的 world-size-4 NCCL 通信组,确认 rank 03 各连接 3 个 Gloo peer,并要求两个 Ray 节点各有 2 个 WorkerDict 进程。脚本还要求 NCCL 使用 Socket 通信,与关闭 RDMA 的任务配置一致。检查结果写入:

language-text
/mnt/verl-reproduction/runs/verl/two-worker-grpo/${RUN_ID}/actor-placement.json

actor-placement.json 中的 status 应为 passedactors 应有 4 条 rank 记录,node_counts 应记录两个不同节点且每个值为 2。这项检查确认一个 world size 为 4 的 ActorRollout group 使用了两台 Worker;本教程不为不同训练阶段指定独立的节点位置。

Step 9 检查训练指标和 checkpoint 分片

任务日志和共享存储应同时满足以下条件:

  • 日志包含 [validate_config] All configuration checks passed successfully!Total training steps: 2
  • TensorBoard 中的 training/global_step 到达 step 2actor/pg_lossactor/grad_normactor/lrcritic/score/mean 在 step 2 均为有限数值;
  • 两份 rollout 文件各包含 8 条记录;
  • fsdp_config.json 记录 FSDP_version=2world_size=4
  • modeloptimextra_state 三类 checkpoint 均包含 rank 03 的非空文件;
  • 日志最后出现 Verification passed: .../verification.json,输出目录包含非空 SUCCESS

任务结束后,在挂载同一共享存储的 AICoder 或开发机中,把 RUN_ID 替换为本次实际值:

language-bash
set -euo pipefail

export RUN_ID=verl-two-worker-grpo-exp-001
export OUTPUT_DIR="/mnt/rlinf-reproduction/runs/verl/two-worker-grpo/$RUN_ID"
export ACTOR_DIR="$OUTPUT_DIR/checkpoints/global_step_2/actor"

test -s "$OUTPUT_DIR/SUCCESS"
test -s "$OUTPUT_DIR/preflight.json"
test -s "$OUTPUT_DIR/actor-placement.json"
test -s "$OUTPUT_DIR/verification.json"
grep -F '"status": "passed"' "$OUTPUT_DIR/verification.json"
grep -F '"world_size": 4' "$ACTOR_DIR/fsdp_config.json"

for kind in model optim extra_state; do
  for rank in 0 1 2 3; do
    test -s "$ACTOR_DIR/${kind}_world_size_4_rank_${rank}.pt"
  done
done

test "$(wc -l < "$OUTPUT_DIR/rollouts/1.jsonl")" -eq 8
test "$(wc -l < "$OUTPUT_DIR/rollouts/2.jsonl")" -eq 8
find "$OUTPUT_DIR/tensorboard" -type f -name 'events.out.tfevents.*' -size +0
du -sh "$OUTPUT_DIR/checkpoints/global_step_2"

这些检查全部通过,说明本次运行完成了目标 step,并把 4 个 rank 的 checkpoint 分片写入共享存储且可以回读。两步结果用于确认训练链路,不用于判断模型是否收敛。

扩大 Worker 或 GPU 数量

本教程的启动脚本和完成检查固定为 2 个 Worker、每个 Worker 2 张 GPU。需要使用其他拓扑时:

  1. 复制启动脚本和校验脚本并创建新的版本化文件。
  2. 同时修改 EXPECTED_NODESGPUS_PER_NODEWORLD_SIZEtrainer.nnodestrainer.n_gpus_per_node
  3. 重新检查 batch、rollout 数量和 world size 的整除关系,以及每张 GPU 的显存需求。
  4. 根据跨节点通信量和目标资源决定是否启用 RDMA;网络配置改变时,同时更新 RuntimeEnv 和完成检查。
  5. 为新拓扑使用新的 RUN_ID,并重新验证 Ray actor 放置、TensorBoard 指标和全部 checkpoint rank。

增加 Worker 或 GPU 不会自动提高这个 1.5B 模型的训练效率。请根据模型规模、batch 和通信开销选择资源,并把新拓扑作为独立运行进行检查。

处理本场景特有故障

Driver 找不到 Ray 集群

确认任务选择了 Ray,并且 RuntimeEnv 和任务环境变量中没有手动设置 RAY_ADDRESS、Head 地址或端口。启动命令不应包含 ray startray job submit

预检只发现一个节点或少于四张 GPU

在任务详情中确认 Worker 数量为 2,每个 Worker 的规格包含 2 张 GPU。修改资源后使用新的 RUN_ID 重跑;不要通过降低脚本中的期望值让错误拓扑通过。

第二个 Worker 找不到代码、模型或数据

确认同一共享存储卷挂载到两个 Worker 的 /mnt/verl-reproduction。RuntimeEnv 中的 working_dirPYTHONPATH、任务环境变量、TensorBoard 路径与启动命令都应使用这个路径。

无法确认四个 ActorRollout rank

查看本次输出目录中的 training.log,确认 4 个 WorkerDict 进程是否都完成 world-size-4 NCCL 初始化、各自连接 3 个 Gloo peer,并使用 Socket 通信。修复最先出现的资源或通信错误后,使用新的 RUN_ID 完整重跑。

日志出现 NCCL 通信错误

确认任务的 RDMA 配置保持关闭,RuntimeEnv 和两个 Ray Worker 中的 NCCL_IB_DISABLE 都为 1。如果要改用 RDMA,需要按新网络配置重新验证镜像、资源和训练,不要只切换页面开关。

任务结束但没有 SUCCESS

查看同一 RUN_ID 目录中的 FAILEDdriver.logtraining.logverification.jsonSUCCESS 只会在两个节点、四个 rank、两个训练 step、有限指标、两份 rollout 和 4 个 rank 的 checkpoint 分片全部通过后生成;修复最先失败的检查,再使用新的 RUN_ID 重跑。