文档首页/ 魔坊(ModelArts)模型训推平台/ 模型训练/ 分布式模型训练/ 示例:创建DDP分布式训练(PyTorch+NPU)
更新时间:2026-07-15 GMT+08:00
分享

示例:创建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控制台训练作业创建分布式作业的过程。

  1. 设置基本信息。设置作业名称,如job-gpu-ddp-example。
    图1 训练基本信息
  2. 训练配置
    • 选择预置镜像,如: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 训练配置
  3. 资源配置
    • 选择有NPU卡的资源池,本案例选择有NPU资源的专属资源池。
    • 规格类型按需选择,示例采用预置规格,如1*Snt9b2 | 24vCPUs | 192GiB,具体规格以环境为准
    • 实例数大于1时,ModelArts默认为分布式训练。本案例设置实例数为2。

    其他参数按需填写。

    图3 资源配置
  4. 更多配置。勾选永久保存日志,日志路径为OBS上定义的输出目录output_dir。这一步为可选操作,也可以不勾选,其他配置参数保持默认即可。
    图4 更多配置

查看执行过程和结果

  1. 在日志页签,用户日志可看到train epoch和loss信息,表示此时训练正在进行中。
    图5 查看训练日志
  2. 事件页签观察训练作业运行成功字样,且作业状态为已完成,表示训练完成。
    图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

相关文档