示例:创建DDP分布式训练(PyTorch+NPU)
在深度学习领域,使用PyTorch的DistributedDataParallel(DDP)功能在Ascend加速卡上进行分布式训练是一种高效的方式。然而,如何通过自定义镜像和自定义启动命令来实现这一目标,是用户在实际操作中可能遇到的挑战。
针对这一问题,本文提供了一种解决方案:通过配置训练作业的自定义镜像,并结合自定义启动命令,用户可以轻松实现PyTorch DDP在Ascend加速卡上的训练任务。这种方法不仅能够满足用户的个性化需求,还能够灵活适配不同的训练场景。
通过本文的介绍,用户可以掌握如何利用自定义镜像和启动命令来优化PyTorch DDP训练流程,从而在Ascend加速卡上实现高效的分布式训练。
前提条件
需要有Ascend加速卡资源池。
准备工作
在创建训练任务之前,需要先准备数据集、准备训练代码。并将上传数据集和训练文件至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'后可直接进行训练作业,无需加载数据。
准备训练代码
run_torch_ddp_npu.sh
其中,PYTHON_SCRIPT参数中的/example_dir需替换为torch_ddp.py脚本实际所在路径。
启动脚本中设置plog生成后存放在“/home/ma-user/modelarts/log/modelarts-job-{id}/worker-{index}/”目录,而“/home/ma-user/modelarts/log/”目录下的“*.log”文件将会被自动上传至ModelArts训练作业的日志目录(OBS)。如果本地相应目录没有生成大小>0的日志文件,则对应的父级目录也不会上传。因此,PyTorch NPU的plog日志是按worker存储的,而不是按rank id存储的(这是区别于MindSpore的)。目前,PyTorch NPU并不依赖rank table file。
#!/bin/bash
# load env variables
source /usr/local/Ascend/ascend-toolkit/set_env.sh
# MA preset envs
MASTER_HOST="$VC_WORKER_HOSTS"
MASTER_ADDR="${VC_WORKER_HOSTS%%,*}"
NNODES="$MA_NUM_HOSTS"
NODE_RANK="$VC_TASK_INDEX"
# also indicates NPU per node
NGPUS_PER_NODE="$MA_NUM_GPUS"
# self-define, it can be changed to >=10000 port
MASTER_PORT="38888"
# replace ${MA_JOB_DIR}/example_dir/torch_ddp.py to the actual training script
PYTHON_SCRIPT=${MA_JOB_DIR}/example_dir/torch_ddp.py
PYTHON_ARGS=""
export HCCL_WHITELIST_DISABLE=1
# set npu plog env
ma_vj_name=`echo ${MA_VJ_NAME} | sed 's:ma-job:modelarts-job:g'`
task_name="worker-${VC_TASK_INDEX}"
task_plog_path=${MA_LOG_DIR}/${ma_vj_name}/${task_name}
mkdir -p ${task_plog_path}
export ASCEND_PROCESS_LOG_PATH=${task_plog_path}
echo "plog path: ${ASCEND_PROCESS_LOG_PATH}"
# set hccl timeout time in seconds
export HCCL_CONNECT_TIMEOUT=1800
# use python from current environment
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
# NPU 适配:尝试导入 torch_npu
try:
import torch_npu
except ImportError:
pass
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)
# 使用 AvgPool2d 替代 AdaptiveAvgPool2d,避免 NPU 上 PadV3 算子兼容性问题
# conv5 输出特征图为 4x4,所以 kernel_size=4 等价于 AdaptiveAvgPool2d((1,1))
self.avg_pool = nn.AvgPool2d(kernel_size=4)
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}')
# 设置当前进程的默认 NPU 设备,确保多卡场景下每个进程绑定到正确的卡
# 不设置时默认设备始终是 npu:0,会导致 local_rank>=1 的进程算子执行失败
torch.npu.set_device(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
tr_loader = DataLoader(tr_set, batch_size=batch_per_gpu, shuffle=False)
### 分布式改造,构建DDP分布式数据sampler,确保不同进程加载到不同的数据 ###
tr_sampler = DistributedSampler(tr_set, num_replicas=world_size, rank=rank)
tr_loader = DataLoader(tr_set, batch_size=batch_per_gpu, sampler=tr_sampler, shuffle=False, drop_last=True)
### 分布式改造,构建DDP分布式数据sampler,确保不同进程加载到不同的数据 ###
val_loader = DataLoader(val_set, batch_size=batch_per_gpu, shuffle=False)
lr = args.lr * 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 训练基本信息
- 训练配置。
- 选择预置镜像,如:2.7.1-cann_8.5.1-py_3.12-hce_2.0.2512-aarch64-snt9b。
- 启动命令如下,其中/DDP-test替换为实际代码所在桶的文件夹。
bash ${MA_JOB_DIR}/DDP-test/run_torch_ddp_npu.sh - 代码目录:选择分布式训练脚本所在桶的文件夹。
其他参数采用默认值。
图2 训练配置
- 资源配置
- 选择有NPU卡的资源池,本案例选择有NPU资源的专属资源池。
- 规格类型按需选择,示例采用预置规格,如1*Snt9b2 | 24vCPUs | 192GiB,具体规格以环境为准
- 实例数大于1时,ModelArts默认为分布式训练。本案例设置实例数为2。
其他参数按需填写。
图3 资源配置
- 更多配置。勾选永久保存日志,日志路径为OBS上定义的输出目录output_dir。这一步为可选操作,也可以不勾选,其他配置参数保持默认即可。 图4 更多配置
查看执行过程和结果
- 在日志页签,用户日志可看到train epoch和loss信息,表示此时训练正在进行中。 图5 查看训练日志
- 事件页签观察训练作业运行成功字样,且作业状态为已完成,表示训练完成。 图6 查看训练结果事件
(可选)启用RankTable动态路由
如果训练作业需要使用RankTable动态路由算法进行网络加速,则可以联系技术支持开启集群的cabinet调度权限。详细使用指导请参见训练作业动态路由加速。
代码示例
训练作业的启动脚本示例如下。
启动脚本中设置plog生成后存放在“/home/ma-user/modelarts/log/modelarts-job-{id}/worker-{index}/”目录,而“/home/ma-user/modelarts/log/”目录下的“*.log”文件将会被自动上传至ModelArts训练作业的日志目录(OBS)。如果本地相应目录没有生成大小>0的日志文件,则对应的父级目录也不会上传。因此,PyTorch NPU的plog日志是按worker存储的,而不是按rank id存储的(这是区别于MindSpore的)。目前,PyTorch NPU并不依赖rank table file。
#!/bin/bash
# load env variables
source /usr/local/Ascend/ascend-toolkit/set_env.sh
# MA preset envs
MASTER_HOST="$VC_WORKER_HOSTS"
MASTER_ADDR="${VC_WORKER_HOSTS%%,*}"
NNODES="$MA_NUM_HOSTS"
NODE_RANK="$VC_TASK_INDEX"
# also indicates NPU per node
NGPUS_PER_NODE="$MA_NUM_GPUS"
# self-define, it can be changed to >=10000 port
MASTER_PORT="38888"
# replace ${MA_JOB_DIR}/code/torch_ddp.py to the actual training script
PYTHON_SCRIPT=${MA_JOB_DIR}/code/torch_ddp.py
PYTHON_ARGS=""
export HCCL_WHITELIST_DISABLE=1
# set npu plog env
ma_vj_name=`echo ${MA_VJ_NAME} | sed 's:ma-job:modelarts-job:g'`
task_name="worker-${VC_TASK_INDEX}"
task_plog_path=${MA_LOG_DIR}/${ma_vj_name}/${task_name}
mkdir -p ${task_plog_path}
export ASCEND_PROCESS_LOG_PATH=${task_plog_path}
echo "plog path: ${ASCEND_PROCESS_LOG_PATH}"
# set hccl timeout time in seconds
export HCCL_CONNECT_TIMEOUT=1800
# replace ${ANACONDA_DIR}/envs/${ENV_NAME}/bin/python to the actual python
CMD="${ANACONDA_DIR}/envs/${ENV_NAME}/bin/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