Skip to main content

Multi-Node Container Training

A single srun launches your container on every node of the allocation. Getting the ranks to talk to each other efficiently is the part that needs care: shared storage, PMIx, and NCCL over the cluster's RDMA fabric.


Prerequisites

RequirementWhyHow to satisfy it
Shared storage on every nodeThe image artefact, the code, the dataset and the checkpoints must be identical and reachable from every rankMount a PFS, SFS or Weka volume; it appears at the same path on the login node and all workers
The image as a squash fileAvoids N nodes pulling the same image and hitting rate limits mid-runenroot import once to shared storage
PMIx for process launchSlurm's default MPI plugin may not be enabled on your regionPass --mpi=pmix on srun
RDMA devices mounted into the container/dev/infiniband is not inside the container by default--container-mounts=…,/dev/infiniband:/dev/infiniband
Nodes in one partition and healthyA drained or unregistered node silently shrinks your allocationCheck sinfo and the Nodes tab
Node-local paths break multi-node jobs

Anything under /tmp, /dev/shm or a node-local scratch path is private to one node. A container name created on slinky-0 does not exist on slinky-1, and a checkpoint written to /tmp on rank 0 is invisible to every other rank. Multi-node means shared storage, without exception.


How Ranks Are Launched

For a container job, srun starts one container per task and Slurm's environment tells each rank who it is.

#SBATCH --nodes=4
#SBATCH --ntasks-per-node=1 # one container per node…
#SBATCH --gpus-per-node=8 # …that owns all 8 GPUs on it

This is the shape to use with torchrun, DeepSpeed or any launcher that spawns its own per-GPU processes: one container per node, the launcher fans out inside it.

For MPI ranks, use one task per GPU instead and let PMIx wire them up:

#SBATCH --nodes=4
#SBATCH --ntasks-per-node=8
#SBATCH --gpus-per-task=1
srun --mpi=pmix --container-image=/pfs/images/hpc.sqsh ./my_mpi_app

Variables every rank gets

VariableMeaning
SLURM_JOB_IDThe job's ID — use it in output paths
SLURM_JOB_NUM_NODESNumber of nodes in the allocation
SLURM_NODEIDThis node's index, 0-based
SLURM_PROCIDThis task's global rank
SLURM_LOCALIDThis task's rank within its node
SLURM_NTASKSTotal tasks — the world size when one task per rank
SLURM_JOB_NODELISTCompact node list, e.g. slinky-[0-3]

Derive the rendezvous host from the node list rather than hard-coding it:

export MASTER_ADDR=$(scontrol show hostnames "$SLURM_JOB_NODELIST" | head -n1)
export MASTER_PORT=29500

NCCL Environment

The platform does not set any NCCL_* or UCX_* variables for your jobs — the defaults come from your image. On a GPU cluster with an RDMA fabric you will normally want to be explicit.

Find the fabric device names first

The HCA names differ by hardware, and they must be identical on every participating node.

srun --nodes=1 --container-image=/pfs/images/hpc.sqsh \
--container-mounts=/dev/infiniband:/dev/infiniband \
bash -lc 'ls /sys/class/infiniband; ibdev2netdev 2>/dev/null'

Run it on more than one node to confirm they agree:

srun --nodes=4 --ntasks-per-node=1 bash -lc 'echo "$(hostname): $(ls /sys/class/infiniband | tr "\n" " ")"'

A working starting point

Adapt the device lists to what the commands above reported:

export NCCL_DEBUG=WARN                  # INFO while debugging, WARN in production
export NCCL_IB_DISABLE=0 # use InfiniBand
export NCCL_NET=IB
export NCCL_IB_HCA=mlx5_0:1,mlx5_3:1,mlx5_4:1,mlx5_5:1,mlx5_6:1,mlx5_9:1,mlx5_10:1,mlx5_11:1
export UCX_NET_DEVICES=mlx5_0:1,mlx5_3:1,mlx5_4:1,mlx5_5:1,mlx5_6:1,mlx5_9:1,mlx5_10:1,mlx5_11:1
export UCX_TLS=rc,cuda
export NCCL_SOCKET_IFNAME=eth0 # interface used for the initial handshake
VariableWhat it controls
NCCL_DEBUGVerbosity. INFO prints the transport NCCL chose — the fastest way to confirm IB is actually in use
NCCL_IB_DISABLE0 to allow IB, 1 to force TCP. Set to 1 only to prove a fabric problem, never for real training
NCCL_IB_HCAWhich HCAs and ports NCCL may use. Wrong names here are the most common cause of a hang
NCCL_SOCKET_IFNAMEInterface for NCCL's bootstrap. Must be a real interface present on every node
UCX_NET_DEVICES, UCX_TLSThe equivalent selection for UCX-based transports (OpenMPI)
NCCL_DEBUG=INFO on a large job is very verbose

Every rank writes its topology decisions. Use it to diagnose, then drop back to WARN. On a 32-rank job the difference is megabytes of log per run.

Confirm IB is really being used

With NCCL_DEBUG=INFO, look for lines naming [send] via NET/IB/…. If you see NET/Socket, NCCL fell back to TCP — your throughput will be a fraction of what the hardware can do. Check NCCL_IB_HCA names and that /dev/infiniband is mounted into the container.


Full Example: PyTorch on 4 Nodes

32 GPUs total, one container per node, torchrun fanning out to 8 local ranks.

#!/bin/bash
#SBATCH --job-name=llm-pretrain
#SBATCH --partition=all
#SBATCH --nodes=4
#SBATCH --ntasks-per-node=1
#SBATCH --gpus-per-node=8
#SBATCH --cpus-per-task=180
#SBATCH --exclusive
#SBATCH --time=24:00:00
#SBATCH --output=/pfs/logs/%x-%j.out
#SBATCH --error=/pfs/logs/%x-%j.err

set -euo pipefail

IMAGE=/pfs/images/pytorch-25.09-py3.sqsh
WORKDIR=/pfs/project
CKPT=/pfs/checkpoints/${SLURM_JOB_NAME}-${SLURM_JOB_ID}
mkdir -p "$CKPT" /pfs/logs

# Rendezvous
export MASTER_ADDR=$(scontrol show hostnames "$SLURM_JOB_NODELIST" | head -n1)
export MASTER_PORT=29500

# Fabric
export NCCL_DEBUG=WARN
export NCCL_IB_DISABLE=0
export NCCL_NET=IB
export NCCL_SOCKET_IFNAME=eth0

echo "Job $SLURM_JOB_ID on $SLURM_JOB_NUM_NODES nodes: $SLURM_JOB_NODELIST"
echo "Rendezvous at $MASTER_ADDR:$MASTER_PORT"

srun --mpi=pmix \
--container-image="$IMAGE" \
--container-name=pretrain \
--container-mounts=/pfs:/pfs,/dev/infiniband:/dev/infiniband \
--container-workdir="$WORKDIR" \
bash -lc '
torchrun \
--nnodes=$SLURM_JOB_NUM_NODES \
--node_rank=$SLURM_NODEID \
--nproc_per_node=8 \
--master_addr=$MASTER_ADDR \
--master_port=$MASTER_PORT \
train.py \
--data /pfs/datasets/corpus \
--checkpoint-dir '"$CKPT"' \
--save-every 500
'

Submit and follow it:

sbatch llm-pretrain.sh
squeue -u $USER
tail -f /pfs/logs/llm-pretrain-<jobid>.out
Checkpoint often on long runs

--save-every 500 is not decoration. A node failure or preemption ends the job; the cluster recovers, but your progress is only as good as your last checkpoint on shared storage. See High availability.

Why --exclusive

--exclusive gives the job whole nodes, so no other job shares the CPU, memory bandwidth or fabric. For large distributed training that usually pays for itself in consistency — a co-tenant saturating the NIC shows up as mysterious step-time variance.


Validate the Fabric First

Before committing a 24-hour run, prove the collectives work at the scale you plan to use. A NCCL all-reduce benchmark is the standard check.

#!/bin/bash
#SBATCH --job-name=nccl-check
#SBATCH --partition=all
#SBATCH --nodes=4
#SBATCH --ntasks-per-node=8
#SBATCH --gpus-per-node=8
#SBATCH --cpus-per-task=20
#SBATCH --exclusive
#SBATCH --time=00:30:00
#SBATCH --output=/pfs/logs/nccl-%j.out

IMAGE=nvcr.io/nvidia/hpc-benchmarks:25.09

export NCCL_DEBUG=INFO
export NCCL_IB_DISABLE=0
export NCCL_NET=IB
export NCCL_SOCKET_IFNAME=eth0

srun --mpi=pmix \
--container-image="$IMAGE" \
--container-name=nccl-bench \
--container-mounts=/dev/infiniband:/dev/infiniband,/pfs:/pfs \
--container-workdir=/workspace \
bash -lc 'all_reduce_perf_mpi -b 8 -e 8G -f 2 -g 1'

Read the busbw column at the largest message size and compare it against the expected bandwidth for your interconnect. A result far below expectation, or a run that hangs at initialisation, points at NCCL_IB_HCA names, a missing /dev/infiniband mount, or one unhealthy node — check the Nodes tab for XID errors.

Also prove shared storage is genuinely shared

A surprising number of multi-node failures are a storage mount that is not what you assumed:

echo "written from $(hostname) at $(date)" > /pfs/shared-check.txt
srun --nodes=4 --ntasks-per-node=1 bash -lc 'echo "$(hostname): $(cat /pfs/shared-check.txt)"'

Every node must print the same line. If one cannot read the file, the volume is not mounted there — check the Volumes tab.


Last updated on September 10, 2026.