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

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

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

    其他参数按需填写。

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

查看执行过程和结果

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

相关文档