已经是最新一篇文章了!
已经是最后一篇文章了!

ai · 2025 · 技术研究笔记

Ray 分布式训练实践

从两个 CPU worker 开始,看清 Ray Train 在分发什么

多机多卡并行训练框架

研究版 · 4 个章节 · 约 4 分钟 · 更新于 2025-01-17

版本 A · 中文
SCROLL

分布式训练最容易产生的误解,是把 worker 数量从 1 改成 4,训练就会快四倍。真实结果取决于计算、通信、数据读取和检查点写入之间的比例。模型较小时,同步梯度的成本可能比计算还高。

Ray Train 的价值不是发明了新的并行算法。它负责启动 worker、分配 CPU 或 GPU、建立 PyTorch 分布式环境,并收集指标与检查点。Worker 是执行训练函数的独立进程。

先在本机跑两个 CPU worker

下面的例子拟合 y = 2x + 1。任务本身很小,因此不会通过分布式获得速度收益;它用来检查 worker、数据切分和指标上报是否连通。

python -m venv .venv
source .venv/bin/activate
python -m pip install "ray[train]" torch
import torch
from torch import nn
from torch.utils.data import DataLoader, TensorDataset

import ray.train
import ray.train.torch
from ray.train import ScalingConfig
from ray.train.torch import TorchTrainer


def train_loop(config: dict[str, float | int]) -> None:
    """在每个 Ray worker 中执行线性回归训练。"""

    x = torch.linspace(-1, 1, 200).reshape(-1, 1)
    y = 2 * x + 1
    dataset = TensorDataset(x, y)
    loader = DataLoader(
        dataset,
        batch_size=int(config["batch_size"]),
        shuffle=True,
    )

    model = nn.Linear(1, 1)
    # prepare_model 会把模型移到当前设备,并包装为分布式模型。
    model = ray.train.torch.prepare_model(model)
    # prepare_data_loader 为每个 worker 分配不同的数据切片。
    loader = ray.train.torch.prepare_data_loader(loader)

    optimizer = torch.optim.SGD(model.parameters(), lr=float(config["lr"]))
    loss_fn = nn.MSELoss()

    for epoch in range(int(config["epochs"])):
        model.train()
        total_loss = 0.0
        batch_count = 0

        for features, labels in loader:
            optimizer.zero_grad()
            prediction = model(features)
            loss = loss_fn(prediction, labels)
            loss.backward()
            optimizer.step()
            total_loss += loss.item()
            batch_count += 1

        mean_loss = total_loss / max(batch_count, 1)
        ray.train.report({"epoch": epoch, "loss": mean_loss})


trainer = TorchTrainer(
    train_loop_per_worker=train_loop,
    train_loop_config={"lr": 0.05, "epochs": 20, "batch_size": 16},
    scaling_config=ScalingConfig(num_workers=2, use_gpu=False),
)

result = trainer.fit()
print(f"最后一次上报的 loss: {result.metrics['loss']:.6f}")

两个 worker 都会执行 train_loop。PyTorch 的分布式数据并行会在反向传播时同步梯度,使各进程的模型参数保持一致。

换成 GPU 时要改哪些

单 worker 使用一张 GPU 时,只需先改资源声明:

scaling_config = ScalingConfig(
    num_workers=4,
    use_gpu=True,
)

Ray 会为每个 worker 分配一张 GPU,并设置对应的环境变量。GPU 默认使用 NCCL 作为通信后端;NCCL 是 NVIDIA 用于多 GPU 集合通信的库。CPU 训练默认使用 Gloo。

多机环境还要确认三件事:

  • 所有节点能访问同一份代码、数据或对象存储。
  • 节点之间的训练网卡选择正确,NCCL 没有走到管理网或慢速接口。
  • 检查点写入共享存储,不是某个 worker 的临时本地目录。

判断扩容是否有效

不要只记录整个任务的墙钟时间。墙钟时间是从任务开始到结束的实际经过时间。至少同时观察:

  • 每秒处理的样本或 Token 数。
  • GPU 利用率和显存使用。
  • 数据加载等待时间。
  • 梯度同步耗时。
  • 检查点写入频率和耗时。

先分别跑 1、2、4 个 worker,并保持全局 batch size 可比。全局 batch size 是所有 worker 在一次参数更新前处理的样本总数。如果 worker 数增加时它也成倍增加,你比较的就不只是并行度,训练语义也变了。

Ray 官方的 DreamBooth 示例中,1、2、4 个 GPU 的训练时间分别约为 802、488 和 313 秒。这个数据说明扩容带来了收益,也说明四张卡不会自动得到四倍速度。它只适用于该示例的模型、硬件和批量设置,不能当作你的集群预测值。

参考资料

版权声明: 如无特别声明,本文版权归 sshipanoo 所有,转载请注明本文链接。

(采用 CC BY-NC-SA 4.0 许可协议进行授权)

本文标题:Ray 分布式训练实践

本文链接:https://www.sshipanoo.com/blog/ai/ray分布式训练/