ai · 2025 · 技术研究笔记
Ray 分布式训练实践
从两个 CPU worker 开始,看清 Ray Train 在分发什么
多机多卡并行训练框架
研究版 · 4 个章节 · 约 4 分钟 · 更新于 2025-01-17
分布式训练最容易产生的误解,是把 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]" torchimport 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 秒。这个数据说明扩容带来了收益,也说明四张卡不会自动得到四倍速度。它只适用于该示例的模型、硬件和批量设置,不能当作你的集群预测值。
