Multi-GPU Implementation for PyTorch on Olivia
This is part 2 of the PyTorch on Olivia guide. See Single-GPU Implementation for PyTorch on Olivia for the single-GPU setup.
Learning Outcomes
By the end of this part, you can:
Run the same training workflow on 4 GPUs on one node.
Understand the minimum DDP changes from the single-GPU version.
Validate that distributed training launched correctly.
To scale training across multiple GPUs, we use PyTorch’s Distributed Data Parallel (DDP). The train_ddp.pycode works for both single-node multi-GPU and multi-node configurations. However, it is important to note that, we don´t use train_utils.py and device_utils.py which we discussed earlier in this page Single-GPU Implementation for PyTorch on Olivia , as train_ddp.py is self contained for DDP and implements equivalent logic directly.
Explanation train_ddp.py file
The train_ddp.py script extends the training pipeline across multiple GPUs and nodes using PyTorch’s DistributedDataParallel (DDP). It coordinates distributed environment setup, multi-process data distribution, and synchronized cross-rank metric collection while strictly adhering to PyTorch distributed training best practices.
It relies on launch utilities like torchrun to set process environment variables (RANK, LOCAL_RANK, WORLD_SIZE) and uses the nccl backend for high-speed inter-GPU communications.
Key Implementation Details & PyTorch DDP Best Practices
Process Group Initialization (
ddp_setup):
Reads environment variables dynamically (RANK, LOCAL_RANK, WORLD_SIZE), binds the process to its explicit local CUDA device, and initializes the process group using the NCCL backend.
Rank 0 Synchronization & Race Prevention:
In a distributed setup, managing filesystem access and logging from a single lead process is critical:
Dataset Downloads: To avoid multi-process race conditions on shared cluster storage,
rank 0performs archive extraction and directory setup exclusively while worker ranks pause at adist.barrier()synchronization point indataset_utils.py.Logging & Early Stopping: Epoch metrics and configuration details are logged exclusively from
rank 0to prevent redundantstdoutnoise across nodes. Whenrank 0triggers an early stopping condition, it communicates this decision to all worker ranks viadist.broadcast()so every process exits synchronously.
Data Sharding (
DistributedSampler):
Divides the global batch size evenly across all active GPUs per_gpu_batch_size = global_batch_size // world_size. This allows scaling the total global batch size seamlessly in Slurm job scripts as node counts increase. Additionally, calling train_sampler.set_epoch(epoch) at the start of every epoch guarantees distinct, non-overlapping dataset shuffles across nodes.
Exact Metric Aggregation (
all_reduce_metrics):
Avoids simple, imprecise averaging of local worker accuracies. Instead, each process tracks raw local sums (correct_count, loss_sum, total_samples) and aggregates them across the cluster using dist.all_reduce(op=dist.ReduceOp.SUM). This produces mathematically exact global loss and accuracy metrics regardless of dataset splitting.
Distributed Throughput Measurement:
Calculates epoch execution times using dist.all_reduce(op=dist.ReduceOp.MAX) across all GPUs. Measuring total processed images against the slowest worker rank provides a true reflection of synchronized step duration across the cluster.
Rank-Specific Seeding & Cleanup:
Seeding: Seeds random number generators with a global rank offset
seed + rank, ensuring worker streams apply independent data augmentations while maintaining overall experiment reproducibility.Cleanup: Enforces process group termination
dist.destroy_process_group()inside a finally block. This guarantees clean execution teardown and prevents orphaned CUDA processes from hanging cluster resources when scaling across nodes.
Job Script for Multi-GPU Training
For single-node multi-GPU training, we should use torchrun command with --standalone. The repo you cloned earlier has the job script where you use the NRIS module. If you choose to use container or EESSI stack the job scripts are given below.
1#!/bin/bash
2#SBATCH --job-name=pytorch_multigpu
3#SBATCH --account=<project_number>
4#SBATCH --output=logs/multigpu_%j.out
5#SBATCH --error=logs/multigpu_%j.err
6#SBATCH --time=00:30:00
7#SBATCH --partition=accel # GPU partition
8#SBATCH --nodes=1 # Single compute node
9#SBATCH --ntasks-per-node=1 # One task (process) on the node
10#SBATCH --cpus-per-task=48 # Right-sized CPU allocation for 4-GPU ViT DDP
11#SBATCH --mem=192G # Right-sized RAM for 4-GPU ViT DDP
12#SBATCH --gpus=4 # Request 4 GPU
13
14# Get the absolute path to the project directory.
15PROJECT_DIR=$(cd "${SLURM_SUBMIT_DIR}/.." && pwd)
16
17# Path to container and training script
18CONTAINER_PATH="/cluster/work/support/container/pytorch_nvidia_25.05_arm64.sif"
19
20TRAINING_SCRIPT="${PROJECT_DIR}/scripts/train_ddp.py --model vit --dataset tiny-imagenet --batch-size 1024 --epochs 100 --optimizer adamw --base-lr 0.0003 --target-accuracy 0.95 --patience 2 --seed 42 --num-workers 8 --amp"
21
22
23
24# Check GPU availability inside the container
25echo "Checking GPU availability inside the container..."
26apptainer exec --nv $CONTAINER_PATH python -c 'import torch; print(torch.cuda.is_available()); print(torch.cuda.device_count())'
27
28# Start GPU utilization monitoring in the background
29GPU_LOG_FILE="${PROJECT_DIR}/jobs/logs/multigpu.log"
30echo "Starting GPU utilization monitoring..."
31nvidia-smi --query-gpu=timestamp,index,name,utilization.gpu,utilization.memory,memory.total,memory.used --format=csv -l 5 > $GPU_LOG_FILE &
32NVIDIA_MONITOR_PID=$!
33
34# Run the training script with torchrun inside the container
35apptainer exec --nv $CONTAINER_PATH torchrun --standalone --nnodes=$SLURM_JOB_NUM_NODES --nproc_per_node=$SLURM_GPUS_ON_NODE $TRAINING_SCRIPT
36
37# Stop GPU utilization monitoring specifically by PID
38echo "Stopping GPU utilization monitoring..."
39kill $NVIDIA_MONITOR_PID
1#!/bin/bash
2#SBATCH --job-name=pytorch_multigpu
3#SBATCH --account=<project_number>
4#SBATCH --output=logs/multigpu_%j.out
5#SBATCH --error=logs/multigpu_%j.err
6#SBATCH --time=00:30:00
7#SBATCH --partition=accel # GPU partition
8#SBATCH --nodes=1 # Single compute node
9#SBATCH --ntasks-per-node=1 # One task (process) on the node
10#SBATCH --cpus-per-task=40 # Reserve 40 CPU cores (Right-sized for 4-GPU WideResNet)
11#SBATCH --mem=128G # Request 128 GB RAM (Right-sized for 4-GPU WideResNet)
12#SBATCH --gpus=4 # Request 4 GPU
13
14module load EESSI/2025.06
15module load torchvision/0.27.0-foss-2025b-PyTorch-2.12.0-CUDA-12.9.1
16
17# Get the absolute path to the project directory.
18PROJECT_DIR=$(cd "${SLURM_SUBMIT_DIR}/.." && pwd)
19
20# Path to the training script
21TRAINING_SCRIPT="${PROJECT_DIR}/scripts/train_ddp.py --model wideresnet --dataset cifar100 --batch-size 1024 --epochs 100 --base-lr 0.04 --target-accuracy 0.95 --patience 2 --seed 42 --amp"
22
23# Change working directory to project root
24cd "${PROJECT_DIR}"
25
26# Check GPU availability
27echo "Checking GPU availability inside the container..."
28python -c 'import torch; print(torch.cuda.is_available()); print(torch.cuda.device_count())'
29
30# Start GPU utilization monitoring in the background
31GPU_LOG_FILE="${PROJECT_DIR}/jobs/logs/multigpu.log"
32echo "Starting GPU utilization monitoring..."
33nvidia-smi --query-gpu=timestamp,index,name,utilization.gpu,utilization.memory,memory.total,memory.used --format=csv -l 5 > $GPU_LOG_FILE &
34NVIDIA_MONITOR_PID=$!
35
36# Run the training script with torchrun inside the container
37torchrun --standalone --nnodes=$SLURM_JOB_NUM_NODES --nproc_per_node=$SLURM_GPUS_ON_NODE $TRAINING_SCRIPT
38
39# Stop GPU utilization monitoring specifically by PID
40echo "Stopping GPU utilization monitoring..."
41kill $NVIDIA_MONITOR_PID
Then you can submit and monitor the running job using these commands:
sbatch multigpu.sh
squeue -u $USER
tail -f multigpu_<jobid>.out
Example output:
Epoch 95/100: time=2.000s, train_loss=0.0077, train_acc=0.9997, val_loss=1.0301, val_acc=0.7474, throughput=24572.7 img/s
Epoch 96/100: time=1.986s, train_loss=0.0078, train_acc=0.9997, val_loss=1.0099, val_acc=0.7461, throughput=24753.2 img/s
Epoch 97/100: time=1.966s, train_loss=0.0095, train_acc=0.9995, val_loss=1.0698, val_acc=0.7380, throughput=24994.7 img/s
Epoch 98/100: time=1.975s, train_loss=0.0090, train_acc=0.9996, val_loss=1.0357, val_acc=0.7513, throughput=24891.3 img/s
Epoch 99/100: time=2.012s, train_loss=0.0081, train_acc=0.9997, val_loss=0.9978, val_acc=0.7515, throughput=24435.0 img/s
Epoch 100/100: time=1.984s, train_loss=0.0082, train_acc=0.9997, val_loss=1.0068, val_acc=0.7491, throughput=24777.5 img/s
Training Summary:
Total training time: 201.998 seconds
Throughput: 24332.868 images/second
Total GPUs used: 4
Training completed successfully.
With 4 GPUs and FP16 mixed precision, the throughput increased from ~7367 images/second (single GPU) to ~24,000 images/second—a 3x speedup. This near-linear scaling demonstrates efficient distributed data parallelism, with slight overhead due to inter-GPU gradient communication.
Success criteria for Part 2:
Output includes
Training started with 4 processesFinal summary reports
Total GPUs used: 4Throughput is substantially higher than Part 1
For Part 3 (multi-node), you keep the same train_ddp.py and only change the job launch configuration. See Multi-Node Guide.