示例:创建DDP分布式训练(PyTorch+GPU)
在分布式训练场景中,使用PyTorch的DistributedDataParallel(DDP)功能是实现高效训练的重要方式。为了帮助用户更好地理解和应用这一功能,本文将详细介绍torch.distributed.launch命令启动训练作业PyTorchDDP训练的方法,并提供对应的代码示例。
前提条件
- 需要有GPU加速卡资源池。
准备工作
在创建训练任务之前,需要先准备数据集、准备训练代码。并将上传数据集和训练文件至OBS。
准备数据集
本文以ResNet18在CIFAR10数据集上的图像分类任务为例。
cifar10数据集
示例代码中提供了三种训练数据加载方式。
cifar-10数据集下载链接,单击“CIFAR-10 python version”。
- 尝试基于torchvision获取cifar10数据集。
- 基于数据链接下载数据并解压,放置在指定目录下,训练集和测试集的大小分别为(50000,3,32,32)和(10000,3,32,32)。
- 考虑到下载cifar10数据集较慢,基于torch生成类似cifar10的随机数据集,训练集和测试集的大小分别为(5000,3,32,32)和(1000,3,32,32),标签仍为10类,指定custom_data = 'true'后可直接进行训练作业,无需加载数据。
准备训练代码
torchlaunch.sh
其中,PYTHON_SCRIPT参数中的/DDP-test需替换为torch_ddp.py脚本实际路径。
#!/bin/bash
# 系统默认环境变量,不建议修改
MASTER_HOST="$VC_WORKER_HOSTS"
MASTER_ADDR="${VC_WORKER_HOSTS%%,*}"
MASTER_PORT="6060"
JOB_ID="1234"
NNODES="$MA_NUM_HOSTS"
NODE_RANK="$VC_TASK_INDEX"
NGPUS_PER_NODE="$MA_NUM_GPUS"
# 自定义环境变量,指定python脚本和参数
PYTHON_SCRIPT=${MA_JOB_DIR}/DDP-test/torch_ddp.py
PYTHON_ARGS=""
CMD="python -m torch.distributed.launch \
--nnodes=$NNODES \
--node_rank=$NODE_RANK \
--nproc_per_node=$NGPUS_PER_NODE \
--master_addr $MASTER_ADDR \
--master_port=$MASTER_PORT \
--use_env \
$PYTHON_SCRIPT \
$PYTHON_ARGS
"
echo $CMD
$CMD torch_ddp.py
import datetime
import inspect
import os
import pickle
import random
import logging
import argparse
import numpy as np
from sklearn.metrics import accuracy_score
import torch
from torch import nn, optim
import torch.distributed as dist
from torch.utils.data import TensorDataset, DataLoader
from torch.utils.data.distributed import DistributedSampler
file_dir = os.path.dirname(inspect.getframeinfo(inspect.currentframe()).filename)
def load_pickle_data(path):
with open(path, 'rb') as file:
data = pickle.load(file, encoding='bytes')
return data
def _load_data(file_path):
raw_data = load_pickle_data(file_path)
labels = raw_data[b'labels']
data = raw_data[b'data']
filenames = raw_data[b'filenames']
data = data.reshape(10000, 3, 32, 32) / 255
return data, labels, filenames
def load_cifar_data(root_path):
train_root_path = os.path.join(root_path, 'cifar-10-batches-py/data_batch_')
train_data_record = []
train_labels = []
train_filenames = []
for i in range(1, 6):
train_file_path = train_root_path + str(i)
data, labels, filenames = _load_data(train_file_path)
train_data_record.append(data)
train_labels += labels
train_filenames += filenames
train_data = np.concatenate(train_data_record, axis=0)
train_labels = np.array(train_labels)
val_file_path = os.path.join(root_path, 'cifar-10-batches-py/test_batch')
val_data, val_labels, val_filenames = _load_data(val_file_path)
val_labels = np.array(val_labels)
tr_data = torch.from_numpy(train_data).float()
tr_labels = torch.from_numpy(train_labels).long()
val_data = torch.from_numpy(val_data).float()
val_labels = torch.from_numpy(val_labels).long()
return tr_data, tr_labels, val_data, val_labels
def get_data(root_path, custom_data=False):
if custom_data:
train_samples, test_samples, img_size = 5000, 1000, 32
tr_label = [1] * int(train_samples / 2) + [0] * int(train_samples / 2)
val_label = [1] * int(test_samples / 2) + [0] * int(test_samples / 2)
random.seed(2021)
random.shuffle(tr_label)
random.shuffle(val_label)
tr_data, tr_labels = torch.randn((train_samples, 3, img_size, img_size)).float(), torch.tensor(tr_label).long()
val_data, val_labels = torch.randn((test_samples, 3, img_size, img_size)).float(), torch.tensor(
val_label).long()
tr_set = TensorDataset(tr_data, tr_labels)
val_set = TensorDataset(val_data, val_labels)
return tr_set, val_set
elif os.path.exists(os.path.join(root_path, 'cifar-10-batches-py')):
tr_data, tr_labels, val_data, val_labels = load_cifar_data(root_path)
tr_set = TensorDataset(tr_data, tr_labels)
val_set = TensorDataset(val_data, val_labels)
return tr_set, val_set
else:
try:
import torchvision
from torchvision import transforms
tr_set = torchvision.datasets.CIFAR10(root='./data', train=True,
download=True, transform=transforms)
val_set = torchvision.datasets.CIFAR10(root='./data', train=False,
download=True, transform=transforms)
return tr_set, val_set
except Exception as e:
raise Exception(
f"{e}, you can download and unzip cifar-10 dataset manually, "
"the data url is http://www.cs.toronto.edu/~kriz/cifar-10-python.tar.gz")
class Block(nn.Module):
def __init__(self, in_channels, out_channels, stride=1):
super().__init__()
self.residual_function = nn.Sequential(
nn.Conv2d(in_channels, out_channels, kernel_size=3, stride=stride, padding=1, bias=False),
nn.BatchNorm2d(out_channels),
nn.ReLU(inplace=True),
nn.Conv2d(out_channels, out_channels, kernel_size=3, padding=1, bias=False),
nn.BatchNorm2d(out_channels)
)
self.shortcut = nn.Sequential()
if stride != 1 or in_channels != out_channels:
self.shortcut = nn.Sequential(
nn.Conv2d(in_channels, out_channels, kernel_size=1, stride=stride, bias=False),
nn.BatchNorm2d(out_channels)
)
def forward(self, x):
out = self.residual_function(x) + self.shortcut(x)
return nn.ReLU(inplace=True)(out)
class ResNet(nn.Module):
def __init__(self, block, num_classes=10):
super().__init__()
self.conv1 = nn.Sequential(
nn.Conv2d(3, 64, kernel_size=3, padding=1, bias=False),
nn.BatchNorm2d(64),
nn.ReLU(inplace=True))
self.conv2 = self.make_layer(block, 64, 64, 2, 1)
self.conv3 = self.make_layer(block, 64, 128, 2, 2)
self.conv4 = self.make_layer(block, 128, 256, 2, 2)
self.conv5 = self.make_layer(block, 256, 512, 2, 2)
self.avg_pool = nn.AdaptiveAvgPool2d((1, 1))
self.dense_layer = nn.Linear(512, num_classes)
def make_layer(self, block, in_channels, out_channels, num_blocks, stride):
strides = [stride] + [1] * (num_blocks - 1)
layers = []
for stride in strides:
layers.append(block(in_channels, out_channels, stride))
in_channels = out_channels
return nn.Sequential(*layers)
def forward(self, x):
out = self.conv1(x)
out = self.conv2(out)
out = self.conv3(out)
out = self.conv4(out)
out = self.conv5(out)
out = self.avg_pool(out)
out = out.view(out.size(0), -1)
out = self.dense_layer(out)
return out
def setup_seed(seed):
torch.manual_seed(seed)
if torch.cuda.is_available():
torch.cuda.manual_seed_all(seed)
if hasattr(torch, 'npu') and torch.npu.is_available():
torch.npu.manual_seed_all(seed)
np.random.seed(seed)
random.seed(seed)
if torch.cuda.is_available():
torch.backends.cudnn.deterministic = True
def obs_transfer(src_path, dst_path):
import moxing as mox
mox.file.copy_parallel(src_path, dst_path)
logging.info(f"end copy data from {src_path} to {dst_path}")
def main():
seed = datetime.datetime.now().year
setup_seed(seed)
parser = argparse.ArgumentParser(description='Pytorch distribute training',
formatter_class=argparse.ArgumentDefaultsHelpFormatter)
parser.add_argument('--device', default='auto', choices=['auto', 'gpu', 'npu', 'cpu'],
help='device type: auto/gpu/npu/cpu')
parser.add_argument('--lr', default='0.01', help='learning rate')
parser.add_argument('--epochs', default='100', help='training iteration')
parser.add_argument('--init_method', default=None, help='tcp_port')
parser.add_argument('--rank', type=int, default=0, help='index of current task')
parser.add_argument('--world_size', type=int, default=1, help='total number of tasks')
parser.add_argument('--custom_data', default='false')
parser.add_argument('--data_url', type=str, default=os.path.join(file_dir, 'input_dir'))
parser.add_argument('--output_dir', type=str, default=os.path.join(file_dir, 'output_dir'))
args, unknown = parser.parse_known_args()
args.custom_data = args.custom_data == 'true'
args.lr = float(args.lr)
args.epochs = int(args.epochs)
# 自动检测设备类型
if args.device == 'auto':
if hasattr(torch, 'npu') and torch.npu.is_available():
args.device = 'npu'
elif torch.cuda.is_available():
args.device = 'gpu'
else:
args.device = 'cpu'
# 根据设备类型确定 backend 和加速卡数量
if args.device == 'npu':
backend = 'hccl'
accelerators_per_node = torch.npu.device_count()
elif args.device == 'gpu':
backend = 'nccl'
accelerators_per_node = torch.cuda.device_count()
else:
backend = 'gloo'
accelerators_per_node = 1
# 获取 local_rank(torch.distributed.launch --use_env 会注入此环境变量)
local_rank = int(os.environ.get('LOCAL_RANK', 0))
# 确定当前进程使用的设备
if args.device == 'npu':
device = torch.device(f'npu:{local_rank}')
elif args.device == 'gpu':
device = torch.device(f'cuda:{local_rank}')
else:
device = torch.device('cpu')
if args.custom_data:
logging.warning('you are training on custom random dataset, '
'validation accuracy may range from 0.4 to 0.6.')
# 从环境变量获取 rank 和 world_size(--use_env 模式下由 torch.distributed.launch 注入)
# 优先使用环境变量,因为 ModelArts 使用 --use_env 启动,不通过命令行参数传递
rank = int(os.environ.get('RANK', args.rank))
world_size = int(os.environ.get('WORLD_SIZE', args.world_size))
init_method = args.init_method or 'env://'
### 分布式改造,DDP初始化进程,其中init_method, rank和world_size参数均由平台自动入参 ###
dist.init_process_group(init_method=init_method, backend=backend, world_size=world_size, rank=rank)
### 分布式改造,DDP初始化进程,其中init_method, rank和world_size参数均由平台自动入参 ###
tr_set, val_set = get_data(args.data_url, custom_data=args.custom_data)
batch_per_gpu = 128
batch = batch_per_gpu * accelerators_per_node
tr_loader = DataLoader(tr_set, batch_size=batch, shuffle=False)
### 分布式改造,构建DDP分布式数据sampler,确保不同进程加载到不同的数据 ###
tr_sampler = DistributedSampler(tr_set, num_replicas=world_size, rank=rank)
tr_loader = DataLoader(tr_set, batch_size=batch, sampler=tr_sampler, shuffle=False, drop_last=True)
### 分布式改造,构建DDP分布式数据sampler,确保不同进程加载到不同的数据 ###
val_loader = DataLoader(val_set, batch_size=batch, shuffle=False)
lr = args.lr * accelerators_per_node * world_size
max_epoch = args.epochs
model = ResNet(Block).to(device)
### 分布式改造,构建DDP分布式模型 ###
if args.device in ('npu', 'gpu'):
model = nn.parallel.DistributedDataParallel(model, device_ids=[local_rank])
else:
model = nn.parallel.DistributedDataParallel(model)
### 分布式改造,构建DDP分布式模型 ###
optimizer = optim.Adam(model.parameters(), lr=lr)
loss_func = torch.nn.CrossEntropyLoss()
os.makedirs(args.output_dir, exist_ok=True)
for epoch in range(1, max_epoch + 1):
model.train()
train_loss = 0
### 分布式改造,DDP sampler, 基于当前的epoch为其设置随机数,避免加载到重复数据 ###
tr_sampler.set_epoch(epoch)
### 分布式改造,DDP sampler, 基于当前的epoch为其设置随机数,避免加载到重复数据 ###
for step, (tr_x, tr_y) in enumerate(tr_loader):
tr_x, tr_y = tr_x.to(device), tr_y.to(device)
out = model(tr_x)
loss = loss_func(out, tr_y)
optimizer.zero_grad()
loss.backward()
optimizer.step()
train_loss += loss.item()
print('train | epoch: %d | loss: %.4f' % (epoch, train_loss / len(tr_loader)))
val_loss = 0
pred_record = []
real_record = []
model.eval()
with torch.no_grad():
for step, (val_x, val_y) in enumerate(val_loader):
val_x, val_y = val_x.to(device), val_y.to(device)
out = model(val_x)
pred_record += list(np.argmax(out.cpu().numpy(), axis=1))
real_record += list(val_y.cpu().numpy())
val_loss += loss_func(out, val_y).item()
val_accu = accuracy_score(real_record, pred_record)
print('val | epoch: %d | loss: %.4f | accuracy: %.4f' % (epoch, val_loss / len(val_loader), val_accu), '\n')
if rank == 0:
# save ckpt every epoch
torch.save(model.state_dict(), os.path.join(args.output_dir, f'epoch_{epoch}.pth'))
if __name__ == '__main__':
main() 上传数据集和训练文件至OBS
将代码和数据集上传到OBS桶,在ModelArts上运行作业时,从OBS桶中读取数据和代码文件。参考目录结构如下:
{OBS bucket} # OBS对象桶,用户可以自定义名称,例如:modelarts-train-bucket
-{OBS file} # OBS文件夹,自定义名称,例如:DDP-examples
- torch_ddp.py # 训练脚本
- torchlaunch.sh # 启动训练作业的启动脚本
- input_dir # OBS文件夹,用于存放训练数据集,可以自定义名称,此处举例为input_dir
- output_dir # OBS文件夹,用于存放训练输出模型,可以自定义名称,此处举例为output_dir 创建训练作业
本节将介绍通过ModelArts控制台训练作业创建分布式作业的过程。
- 设置基本信息。设置作业名称,如job-gpu-ddp-example。 图1 训练基本信息
- 训练配置。
- 选择预置镜像,如:1.12.1-cuda_10.2-py_3.9.11-ubuntu_22.04-x86_64
- 启动命令如下,其中/DDP-test替换为实际代码所在桶的文件夹。
bash ${MA_JOB_DIR}/DDP-test/torchlaunch.sh - 代码目录:选择分布式训练脚本所在桶的文件夹。
其他参数采用默认值。
图2 训练配置
- 资源配置
- 选择有GPU卡的资源池,本案例选择有GPU资源的专属资源池。
- 规格类型按需选择,示例采用预置规格,如1*pnt004 | 8vCPUs | 32GiB,具体规格以环境为准
- 实例数大于1时,ModelArts默认为分布式训练。本案例设置实例数为2。
其他参数按需填写。
图3 资源配置
- 更多配置。勾选永久保存日志,日志路径为OBS上定义的输出目录output_dir。这一步为可选操作,也可以不勾选。 图4 更多配置
查看执行过程和结果
- 在日志页签,用户日志可看到train epoch和loss信息,表示此时训练正在进行中。 图5 查看训练日志
- 事件页签观察训练作业运行成功字样,且作业状态为已完成,表示训练完成。 图6 查看训练结果事件