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:

  1. Run the same training workflow on 4 GPUs on one node.

  2. Understand the minimum DDP changes from the single-GPU version.

  3. 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

  1. 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.

  1. 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 0 performs archive extraction and directory setup exclusively while worker ranks pause at a dist.barrier() synchronization point in dataset_utils.py.

  • Logging & Early Stopping: Epoch metrics and configuration details are logged exclusively from rank 0 to prevent redundant stdout noise across nodes. When rank 0 triggers an early stopping condition, it communicates this decision to all worker ranks via dist.broadcast() so every process exits synchronously.

  1. 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.

  1. 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.

  1. 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.

  1. 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

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 processes

  • Final summary reports Total GPUs used: 4

  • Throughput 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.