---
title: Multi-Node Container Training
sidebar_label: Multi-node training
---
import { Share2, Cpu, Activity, Layers, AlertTriangle } from 'react-feather';
# 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.
},
{ href: '#how-ranks-are-launched', label: 'How ranks launch', icon: },
{ href: '#nccl-environment', label: 'NCCL environment', icon: },
{ href: '#full-example-pytorch-on-4-nodes', label: 'PyTorch on 4 nodes', icon: },
{ href: '#validate-the-fabric-first', label: 'Validate the fabric', icon: },
]} />
---
## Prerequisites
| Requirement | Why | How to satisfy it |
|-------------|-----|-------------------|
| **Shared storage on every node** | The image artefact, the code, the dataset and the checkpoints must be identical and reachable from every rank | Mount a PFS, SFS or Weka volume; it appears at the same path on the login node and all workers |
| **The image as a squash file** | Avoids N nodes pulling the same image and hitting rate limits mid-run | [`enroot import` once to shared storage](/docs/tir/SlurmCluster/containers/image-cache#import-once-to-a-squash-file) |
| **PMIx for process launch** | Slurm's default MPI plugin may not be enabled on your region | Pass `--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 healthy** | A drained or unregistered node silently shrinks your allocation | Check `sinfo` and the [Nodes tab](/docs/tir/SlurmCluster/manage/nodes) |
:::danger 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.
```bash
#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:
```bash
#SBATCH --nodes=4
#SBATCH --ntasks-per-node=8
#SBATCH --gpus-per-task=1
```
```bash
srun --mpi=pmix --container-image=/pfs/images/hpc.sqsh ./my_mpi_app
```
### Variables every rank gets
| Variable | Meaning |
|----------|---------|
| `SLURM_JOB_ID` | The job's ID — use it in output paths |
| `SLURM_JOB_NUM_NODES` | Number of nodes in the allocation |
| `SLURM_NODEID` | This node's index, `0`-based |
| `SLURM_PROCID` | This task's global rank |
| `SLURM_LOCALID` | This task's rank within its node |
| `SLURM_NTASKS` | Total tasks — the world size when one task per rank |
| `SLURM_JOB_NODELIST` | Compact node list, e.g. `slinky-[0-3]` |
Derive the rendezvous host from the node list rather than hard-coding it:
```bash
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**.
```bash
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:
```bash
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:
```bash
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
```
| Variable | What it controls |
|----------|------------------|
| `NCCL_DEBUG` | Verbosity. `INFO` prints the transport NCCL chose — the fastest way to confirm IB is actually in use |
| `NCCL_IB_DISABLE` | `0` to allow IB, `1` to force TCP. Set to `1` only to prove a fabric problem, never for real training |
| `NCCL_IB_HCA` | Which HCAs and ports NCCL may use. Wrong names here are the most common cause of a hang |
| `NCCL_SOCKET_IFNAME` | Interface for NCCL's bootstrap. Must be a real interface present on every node |
| `UCX_NET_DEVICES`, `UCX_TLS` | The equivalent selection for UCX-based transports (OpenMPI) |
:::warning `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.
:::
:::tip 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.
```bash
#!/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:
```bash
sbatch llm-pretrain.sh
squeue -u $USER
tail -f /pfs/logs/llm-pretrain-.out
```
:::tip 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](/docs/tir/SlurmCluster/getting-started/slurm-cluster-concepts#high-availability-and-what-it-does-not-cover).
:::
### 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.
```bash
#!/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](/docs/tir/SlurmCluster/manage/nodes) 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:
```bash
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](/docs/tir/SlurmCluster/manage/storage).
---
## Related Resources
- [Containers overview](/docs/tir/SlurmCluster/containers/)
- [Run containers in jobs](/docs/tir/SlurmCluster/containers/run-containers)
- [Cache and manage images](/docs/tir/SlurmCluster/containers/image-cache)
- [Troubleshoot container jobs](/docs/tir/SlurmCluster/troubleshoot/containers)
- [Nodes and GPU health](/docs/tir/SlurmCluster/manage/nodes)