Hugging Face 教程:从 PyTorch DDP 到 Accelerate 再到 Trainer 掌握分布式训练
From PyTorch DDP to Accelerate to Trainer, mastery of distributed training with ease
Hugging Face 发布教程,以 MNIST 训练示例讲解三种抽象层级的分布式数据并行训练:PyTorch 原生 torch.distributed 与 DDP、Accelerate 的轻量封装、Transformers 的 Trainer API。
教程用同一个 MNIST 示例逐级展示三种分布式训练写法,读者可以按抽象层级选择适合自己项目的方案。
总体概述
本教程假设您对 PyTorch 以及如何训练一个简单模型有基本的了解。它将通过三个抽象级别递增的方式,展示如何通过称为分布式数据并行(DDP)的过程在多个 GPU 上进行训练:
- 通过
pytorch.distributed模块使用原生 PyTorch DDP - 利用 🤗 Accelerate 对
pytorch.distributed的轻量封装,这也有助于确保代码可以在单个 GPU 和 TPU 上运行,无需更改代码,并且对原始代码的更改也极少 - 利用 🤗 Transformer 的高级 Trainer API,它抽象了所有样板代码,并支持各种设备和分布式场景
什么是“分布式”训练,为什么它很重要?
看看下面一些非常基础的 PyTorch 训练代码,它基于官方 MNIST 示例在 MNIST 上设置并训练一个模型
import torch
import torch.nn as nn
import torch.nn.functional as F
import torch.optim as optim
from torchvision import datasets, transforms
class BasicNet(nn.Module):
def __init__(self):
super().__init__()
self.conv1 = nn.Conv2d(1, 32, 3, 1)
self.conv2 = nn.Conv2d(32, 64, 3, 1)
self.dropout1 = nn.Dropout(0.25)
self.dropout2 = nn.Dropout(0.5)
self.fc1 = nn.Linear(9216, 128)
self.fc2 = nn.Linear(128, 10)
self.act = F.relu
def forward(self, x):
x = self.act(self.conv1(x))
x = self.act(self.conv2(x))
x = F.max_pool2d(x, 2)
x = self.dropout1(x)
x = torch.flatten(x, 1)
x = self.act(self.fc1(x))
x = self.dropout2(x)
x = self.fc2(x)
output = F.log_softmax(x, dim=1)
return output
我们定义训练设备(cuda):
device = "cuda"
构建一些 PyTorch DataLoader:
transform = transforms.Compose([
transforms.ToTensor(),
transforms.Normalize((0.1307), (0.3081))
])
train_dset = datasets.MNIST('data', train=True, download=True, transform=transform)
test_dset = datasets.MNIST('data', train=False, transform=transform)
train_loader = torch.utils.data.DataLoader(train_dset, shuffle=True, batch_size=64)
test_loader = torch.utils.data.DataLoader(test_dset, shuffle=False, batch_size=64)
将模型移动到 CUDA 设备:
model = BasicNet().to(device)
构建一个 PyTorch 优化器:
optimizer = optim.AdamW(model.parameters(), lr=1e-3)
最后创建一个简单的训练和评估循环,对数据集执行一次完整的迭代并计算测试准确率:
model.train()
for batch_idx, (data, target) in enumerate(train_loader):
data, target = data.to(device), target.to(device)
output = model(data)
loss = F.nll_loss(output, target)
loss.backward()
optimizer.step()
optimizer.zero_grad()
model.eval()
correct = 0
with torch.no_grad():
for data, target in test_loader:
data, target = data.to(device), target.to(device)
output = model(data)
pred = output.argmax(dim=1, keepdim=True)
correct += pred.eq(target.view_as(pred)).sum().item()
print(f'Accuracy: {100. * correct / len(test_loader.dataset)}')
通常从这里开始,你可以将所有这些内容放入一个 Python 脚本中,或者在 Jupyter Notebook 中运行它。
然而,如果这些资源可用,你如何让这个脚本在比如两个 GPU 或多台机器上运行,从而通过分布式训练提高训练速度呢?仅仅执行 python myscript.py 只会使用单个 GPU 运行脚本。这就是 torch.distributed 发挥作用的地方
PyTorch 分布式数据并行
顾名思义,torch.distributed 旨在用于分布式设置。这可能包括多节点,即你有若干台机器,每台机器有一个 GPU,或者多 GPU,即单个系统有多个 GPU,或者两者的某种组合。
要将我们上面的代码转换为在分布式设置中工作,必须首先定义一些设置配置,详见DDP 入门教程
首先必须声明一个 setup 和一个 cleanup 函数。这将打开一个处理组,所有计算进程都可以通过它进行通信
注意:在本教程的这一部分,应假设这些内容是在 Python 脚本文件中发送的。稍后将讨论使用 Accelerate 的启动器,它可以消除这种必要性
import os
import torch.distributed as dist
def setup(rank, world_size):
"Sets up the process group and configuration for PyTorch Distributed Data Parallelism"
os.environ["MASTER_ADDR"] = 'localhost'
os.environ["MASTER_PORT"] = "12355"
# Initialize the process group
dist.init_process_group("gloo", rank=rank, world_size=world_size)
def cleanup():
"Cleans up the distributed environment"
dist.destroy_process_group()
最后一块拼图是如何将我的数据和模型发送到另一个 GPU?
这就是 DistributedDataParallel 模块发挥作用的地方。它会将你的模型复制到每个 GPU 上,当调用 loss.backward() 时,会执行反向传播,并且所有这些模型副本上产生的梯度将被平均/归约。这确保每个设备在优化器步骤后具有相同的权重。
下面是我们训练设置的示例,将其重构为一个函数,并具备此能力:
注意:这里的 rank 是当前 GPU 相对于所有其他可用 GPU 的总体排名,这意味着它们的 rank 为
0 -> n-1
from torch.nn.parallel import DistributedDataParallel as DDP
def train(model, rank, world_size):
setup(rank, world_size)
model = model.to(rank)
ddp_model = DDP(model, device_ids=[rank])
optimizer = optim.AdamW(ddp_model.parameters(), lr=1e-3)
# Train for one epoch
ddp_model.train()
for batch_idx, (data, target) in enumerate(train_loader):
data, target = data.to(rank), target.to(rank)
output = model(data)
loss = F.nll_loss(output, target)
loss.backward()
optimizer.step()
optimizer.zero_grad()
cleanup()
优化器需要基于特定设备上的模型来声明(因此是 ddp_model 而不是 model),以便正确计算所有梯度。
最后,要运行脚本,PyTorch 有一个方便的 torchrun 命令行模块可以提供帮助。只需传入它应使用的节点数以及要运行的脚本,你就可以开始了:
torchrun --nproc_per_node=2 --nnodes=1 example_script.py
以上将在单台机器上的两个 GPU 上运行训练脚本,这是仅使用 PyTorch 执行分布式训练的基础。
现在让我们谈谈 Accelerate,这个库旨在让这一过程更加无缝,并帮助遵循一些最佳实践。
🤗 Accelerate
Accelerate 是一个库,旨在让你无需大幅修改代码即可执行我们刚才所做的操作。此外,Accelerate 固有的数据管道还能提升代码性能。
首先,让我们把刚才执行的所有代码封装成一个函数,以便直观地看到差异:
def train_ddp(rank, world_size):
setup(rank, world_size)
# Build DataLoaders
transform = transforms.Compose([
transforms.ToTensor(),
transforms.Normalize((0.1307), (0.3081))
])
train_dset = datasets.MNIST('data', train=True, download=True, transform=transform)
test_dset = datasets.MNIST('data', train=False, transform=transform)
train_loader = torch.utils.data.DataLoader(train_dset, shuffle=True, batch_size=64)
test_loader = torch.utils.data.DataLoader(test_dset, shuffle=False, batch_size=64)
# Build model
model = model.to(rank)
ddp_model = DDP(model, device_ids=[rank])
# Build optimizer
optimizer = optim.AdamW(ddp_model.parameters(), lr=1e-3)
# Train for a single epoch
ddp_model.train()
for batch_idx, (data, target) in enumerate(train_loader):
data, target = data.to(rank), target.to(rank)
output = ddp_model(data)
loss = F.nll_loss(output, target)
loss.backward()
optimizer.step()
optimizer.zero_grad()
# Evaluate
model.eval()
correct = 0
with torch.no_grad():
for data, target in test_loader:
data, target = data.to(rank), target.to(rank)
output = ddp_model(data)
pred = output.argmax(dim=1, keepdim=True)
correct += pred.eq(target.view_as(pred)).sum().item()
print(f'Accuracy: {100. * correct / len(test_loader.dataset)}')
接下来,我们讨论 Accelerate 如何提供帮助。上述代码存在几个问题:
- 这略显低效,因为
n数据加载器是基于每个设备创建并推送的。 - 此代码仅适用于多 GPU,因此若要在单节点或 TPU 上运行,需要特别处理。
Accelerate 通过 Accelerator 类解决了这个问题。通过它,代码在单节点与多节点之间仅需三行代码的差异,如下所示:
def train_ddp_accelerate():
accelerator = Accelerator()
# Build DataLoaders
transform = transforms.Compose([
transforms.ToTensor(),
transforms.Normalize((0.1307), (0.3081))
])
train_dset = datasets.MNIST('data', train=True, download=True, transform=transform)
test_dset = datasets.MNIST('data', train=False, transform=transform)
train_loader = torch.utils.data.DataLoader(train_dset, shuffle=True, batch_size=64)
test_loader = torch.utils.data.DataLoader(test_dset, shuffle=False, batch_size=64)
# Build model
model = BasicNet()
# Build optimizer
optimizer = optim.AdamW(model.parameters(), lr=1e-3)
# Send everything through `accelerator.prepare`
train_loader, test_loader, model, optimizer = accelerator.prepare(
train_loader, test_loader, model, optimizer
)
# Train for a single epoch
model.train()
for batch_idx, (data, target) in enumerate(train_loader):
output = model(data)
loss = F.nll_loss(output, target)
accelerator.backward(loss)
optimizer.step()
optimizer.zero_grad()
# Evaluate
model.eval()
correct = 0
with torch.no_grad():
for data, target in test_loader:
data, target = data.to(device), target.to(device)
output = model(data)
pred = output.argmax(dim=1, keepdim=True)
correct += pred.eq(target.view_as(pred)).sum().item()
print(f'Accuracy: {100. * correct / len(test_loader.dataset)}')
有了这个,你的 PyTorch 训练循环现在可以借助 Accelerator 对象在任何分布式设置上运行。然后,这段代码仍可通过 torchrun CLI 或 Accelerate 自有的 CLI 接口 accelerate launch 启动。
因此,使用 Accelerate 进行分布式训练变得轻而易举,同时尽可能保持 PyTorch 基础代码不变。
之前提到,Accelerate 还能让 DataLoaders 更高效。这是通过自定义 Samplers 实现的,它们可以在训练期间自动将批次的部分发送到不同设备,从而允许一次只知晓一份数据副本,而不是根据配置一次性将四份数据加载到内存中。此外,原始数据集在内存中总共只有一份完整副本。该数据集的子集在所有用于训练的节点之间分割,从而允许在单个实例上训练更大的数据集,而不会导致内存使用量激增。
使用 notebook_launcher
之前提到,你可以直接从 Jupyter Notebook 启动分布式代码。这得益于 Accelerate 的 notebook_launcher 实用工具,它允许基于 Jupyter Notebook 内的代码启动多 GPU 训练。
使用它就像导入启动器一样简单:
from accelerate import notebook_launcher
并传入我们之前声明的训练函数、任何要传递的参数以及要使用的进程数(例如 TPU 上为 8,或两个 GPU 为 2)。上述两个训练函数都可以运行,但请注意,启动一次后,实例需要重启才能再次启动。
notebook_launcher(train_ddp, args=(), num_processes=2)
或者:
notebook_launcher(train_ddp_accelerate, args=(), num_processes=2)
使用 🤗 Trainer
最后,我们来到了最高级别的 API —— Hugging Face Trainer。
它尽可能多地封装了训练过程,同时仍能在分布式系统上训练,而用户无需进行任何操作。
首先,我们需要导入 Trainer:
from transformers import Trainer
然后,我们定义一些 TrainingArguments 来控制所有常规超参数。Trainer 也通过字典工作,因此需要制作自定义的整理函数。
最后,我们子类化 Trainer 并编写自己的 compute_loss。
之后,这段代码也能在分布式设置上运行,而无需编写任何训练代码!
from transformers import Trainer, TrainingArguments
model = BasicNet()
training_args = TrainingArguments(
"basic-trainer",
per_device_train_batch_size=64,
per_device_eval_batch_size=64,
num_train_epochs=1,
evaluation_strategy="epoch",
remove_unused_columns=False
)
def collate_fn(examples):
pixel_values = torch.stack([example[0] for example in examples])
labels = torch.tensor([example[1] for example in examples])
return {"x":pixel_values, "labels":labels}
class MyTrainer(Trainer):
def compute_loss(self, model, inputs, return_outputs=False):
outputs = model(inputs["x"])
target = inputs["labels"]
loss = F.nll_loss(outputs, target)
return (loss, outputs) if return_outputs else loss
trainer = MyTrainer(
model,
training_args,
train_dataset=train_dset,
eval_dataset=test_dset,
data_collator=collate_fn,
)
trainer.train()
***** Running training *****
Num examples = 60000
Num Epochs = 1
Instantaneous batch size per device = 64
Total train batch size (w. parallel, distributed & accumulation) = 64
Gradient Accumulation steps = 1
Total optimization steps = 938
| 轮次 | 训练损失 | 验证损失 |
|---|---|---|
| 1 | 0.875700 | 0.282633 |
与上述使用 notebook_launcher 的示例类似,这里也可以通过将所有内容放入一个训练函数来再次实现:
def train_trainer_ddp():
model = BasicNet()
training_args = TrainingArguments(
"basic-trainer",
per_device_train_batch_size=64,
per_device_eval_batch_size=64,
num_train_epochs=1,
evaluation_strategy="epoch",
remove_unused_columns=False
)
def collate_fn(examples):
pixel_values = torch.stack([example[0] for example in examples])
labels = torch.tensor([example[1] for example in examples])
return {"x":pixel_values, "labels":labels}
class MyTrainer(Trainer):
def compute_loss(self, model, inputs, return_outputs=False):
outputs = model(inputs["x"])
target = inputs["labels"]
loss = F.nll_loss(outputs, target)
return (loss, outputs) if return_outputs else loss
trainer = MyTrainer(
model,
training_args,
train_dataset=train_dset,
eval_dataset=test_dset,
data_collator=collate_fn,
)
trainer.train()
notebook_launcher(train_trainer_ddp, args=(), num_processes=2)
资源
要了解有关 PyTorch 分布式数据并行的更多信息,请查看此处的文档
要了解有关 🤗 Accelerate 的更多信息,请查看此处的文档
要了解有关 🤗 Transformers 的更多信息,请查看此处的文档
来源:Hugging Face:Blog · huggingface.co