Help Center/ ModelArts/ Model Training/ Distributed Model Training/ Example: Creating a DDP Distributed Training Job (PyTorch + NPU)
Updated on 2026-08-27 GMT+08:00

Example: Creating a DDP Distributed Training Job (PyTorch + NPU)

PyTorch's DDP enables efficient distributed training on Ascend accelerator cards. Implementing this with custom images and boot commands can pose challenges.

This guide offers a solution: configure custom images and boot commands to run PyTorch DDP training on Ascend accelerator cards. It supports user-specific needs and adapts to various training scenarios.

With this approach, users can streamline PyTorch DDP training for better performance on Ascend accelerator cards.

Prerequisites

An Ascend accelerator card resource pool is available.

Preparations

Before creating a training task, you need to perform operations in Preparing a Dataset, Preparing Training Code, and Uploading the Dataset and Training Files to OBS.

Preparing a Dataset

This section uses ResNet18 as an example to describe how to perform an image classification task on the CIFAR-10 dataset.

CIFAR-10 dataset

The following provides three methods to load training data.

Click CIFAR-10 python version on the download page to download the CIFAR-10 dataset.

  • Download the CIFAR-10 dataset using torchvision.
  • Download the CIFAR-10 dataset based on the URL and decompress the dataset in a specified directory. The sizes of the training set and test set are (50000, 3, 32, 32) and (10000, 3, 32, 32), respectively.
  • Use Torch to obtain a random dataset similar to CIFAR-10. The sizes of the training set and test set are (5000, 3, 32, 32) and (1000, 3, 32, 32), respectively. The labels are still of 10 types. Set custom_data to true, and the training task can be directly executed without loading data.

Preparing Training Code

run_torch_ddp_npu.sh

Replace /example_dir in the PYTHON_SCRIPT parameter with the actual path of the torch_ddp.py script.

To store the generated plog data, you need to specify the path in the startup script as /home/ma-user/modelarts/log/modelarts-job-{id}/worker-{index}/. The system will automatically upload the *.log file in the /home/ma-user/modelarts/log/ directory to the OBS log directory of your training job. The system will only upload the log files (larger than 0 MB) in the local directory to the corresponding parent directory. Unlike MindSpore, PyTorch NPU plog logs are organized by worker instead of rank ID. PyTorch NPU does not rely on the ranktable 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 adaptation: Try to import 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)
        # Replace AdaptiveAvgPool2d with AvgPool2d to avoid the compatibility issue of the PadV3 operator on the NPUs.
        # The output feature map of conv5 is 4 x 4. Therefore, kernel_size=4 is equivalent to 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)
    # Automatically detect device type
    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'
    # Determine backend and number of accelerator cards based on device type
    if args.device == 'npu':
        backend = 'hccl'
    elif args.device == 'gpu':
        backend = 'nccl'
    else:
        backend = 'gloo'
    # Get local_rank (injected by torch.distributed.launch --use_env)
    local_rank = int(os.environ.get('LOCAL_RANK', 0))
    # Determine the device used by the current process
    if args.device == 'npu':
        device = torch.device(f'npu:{local_rank}')
        # Set the default NPU device for the current process to ensure that each process is bound to the correct NPU in the multi-PU scenario.
        # If this parameter is not set, the default device is always npu:0. As a result, the process operator whose local_rank is greater than or equal to 1 fails to be executed.
        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.')
    # Get rank and world_size from environment variables (injected by torch.distributed.launch in --use_env mode)
    # Prioritize environment variables because ModelArts uses --use_env to launch and does not pass command-line arguments
    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://'
    ### Settings for distributed training. Initialize DistributedDataParallel process. The init_method, rank, and world_size parameters are automatically input by the platform. ###
    dist.init_process_group(init_method=init_method, backend=backend, world_size=world_size, rank=rank)
    ### Settings for distributed training. Initialize DistributedDataParallel process. The init_method, rank, and world_size parameters are automatically input by the platform. ###
    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)
    ### Settings for distributed training. Create a sampler for data distribution to ensure that different processes load different data. ###
    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)
    ### Settings for distributed training. Create a sampler for data distribution to ensure that different processes load different data. ###
    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)
    ### Settings for distributed training. Build a DistributedDataParallel model. ###
    if args.device in ('npu', 'gpu'):
        model = nn.parallel.DistributedDataParallel(model, device_ids=[local_rank])
    else:
        model = nn.parallel.DistributedDataParallel(model)
    ### Settings for distributed training. Build a DistributedDataParallel model. ###
    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
        ### Settings for distributed training. DistributedDataParallel sampler. Random numbers are set for the DistributedDataParallel sampler based on the current epoch number to avoid loading duplicate data. ###
        tr_sampler.set_epoch(epoch)
        ### Settings for distributed training. DistributedDataParallel sampler. Random numbers are set for the DistributedDataParallel sampler based on the current epoch number to avoid loading duplicate data. ###
        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()

Uploading the Dataset and Training Files to OBS

Upload your code and dataset to an OBS bucket. When running a job on ModelArts, the system will read the data and code files directly from the OBS bucket. Refer to the directory structure below:

{OBS bucket}                     # OBS bucket name, which is customizable (e.g., modelarts-train-bucket)
    -{OBS file}                  # OBS folder name, which is customizable (e.g., DDP-examples)
        - torch_ddp.py # Training script
        - torchlaunch.sh         # Launch script to start the training job
        - input_dir              # OBS folder for storing the training dataset. The folder name is customizable (e.g., input_dir).
        - output_dir             # OBS folder for storing trained model outputs. The folder name is customizable (e.g., output_dir).

Creating a Training Job

This section describes how to create a distributed training job via the ModelArts console.

  1. Configure basic information. Set the job name, for example, job-npu-ddp-example.
  2. Finish training configuration.
    • Select Image: Select a preset image, for example, 2.7.1-cann_8.5.1-py_3.12-hce_2.0.2512-aarch64-snt9b.
    • Boot Command: Use the command below (replace /DDP-test with the actual folder in your bucket where the code resides):
      bash ${MA_JOB_DIR}/DDP-test/run_torch_ddp_npu.sh
    • Code Directory: Select the bucket folder containing the distributed training script.

    Retain default settings for other parameters.

  3. Finish resource configuration.
    • Resource Pool: Select a resource pool with NPUs. In this example, select a dedicated resource pool with NPU resources.
    • Specification Type: Select a specification according to your requirements. This example uses a preset specification, such as 1*Snt9b2 | 24vCPUs | 192GiB (actual specifications depend on your environment).
    • Compute Nodes: When set to greater than 1, ModelArts defaults to distributed training. In this example, set this parameter to 2.

    Set other parameters as required.

  4. Finish more configurations. Persistent Log Saving: Check the box to permanently save logs, setting the log path to the output_dir directory defined on OBS. This step is optional and can be skipped. Retain default settings for other parameters.

Viewing the Execution Process and Result

  1. In the Logs tab, view the train epoch and loss information, indicating that the training is in progress.
  2. In the Events tab, check the training job status. If Training job completed is displayed and the job status is Completed, the training is complete.

(Optional) Enabling Ranktable Dynamic Routing

To use ranktable dynamic routing for network acceleration, contact technical support to enable cabinet scheduling permission. For details, see Enabling Dynamic Route Acceleration for Training Jobs.

Code Example

The following shows an example boot script of a training job.

To store the generated plog data, you need to specify the path in the startup script as /home/ma-user/modelarts/log/modelarts-job-{id}/worker-{index}/. The system will automatically upload the *.log file in the /home/ma-user/modelarts/log/ directory to the OBS log directory of your training job. The system will only upload the log files (larger than 0 MB) in the local directory to the corresponding parent directory. Unlike MindSpore, PyTorch NPU plog logs are organized by worker instead of rank ID. PyTorch NPU does not rely on the ranktable 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