Distributed GPU Training with NCCL
The NVIDIA Collective Communications Library (NCCL) is the communication layer that lets multiple GPUs exchange data during distributed jobs – the GPU counterpart to MPI. It is the backbone of multi-GPU AI training.
NCCL is a library, not a process launcher: you do not wire it up directly.
Your framework – almost always PyTorch – uses NCCL internally, and a launcher
such as torchrun starts the processes and assigns ranks. Because NCCL is not a
wire-up mechanism, Fuzzball has no dedicated NCCL multinode.implementation.
Instead you run these workloads with the
generic implementation,
which hands your launcher the per-node rank and head-node address it needs.
A single node with several GPUs needs no multinode block at all. torchrun --standalone launches one process per GPU on the node, and PyTorch uses NCCL to
communicate between them:
version: v4
# Single-node, multi-GPU PyTorch/NCCL run.
#
# A single node with multiple GPUs does NOT need the multinode block: torchrun
# --standalone launches one process per GPU, and PyTorch uses NCCL to
# communicate between them. Fuzzball sets RLIMIT_MEMLOCK to the job's memory
# automatically, which covers NCCL's pinned buffers.
jobs:
torch-allreduce:
image:
uri: docker://nvcr.io/nvidia/pytorch:25.09-py3
env:
- NCCL_DEBUG=INFO
script: |
#!/bin/sh
set -eux
cat > /tmp/allreduce.py <<'PY'
import os, torch
import torch.distributed as dist
dist.init_process_group("nccl")
torch.cuda.set_device(int(os.environ["LOCAL_RANK"]))
t = torch.ones(1, device="cuda") * (dist.get_rank() + 1)
dist.all_reduce(t)
print(f"rank {dist.get_rank()}/{dist.get_world_size()} sum={t.item()}", flush=True)
dist.destroy_process_group()
PY
# --nproc-per-node must match the GPU count requested below.
torchrun --standalone --nproc-per-node=2 /tmp/allreduce.py
resource:
cpu:
cores: 2
affinity: NUMA
memory:
size: 16GB
devices:
nvidia.com/gpu: 2
Set --nproc-per-node to the GPU count. Fuzzball automatically raises the job’s
locked-memory limit (RLIMIT_MEMLOCK) to its memory allocation, which covers
NCCL’s pinned communication buffers.
To span nodes, use the generic implementation. It exposes the per-node rank as
$RANK and the head node’s address as $MULTINODE_NODE_IP – exactly
torchrun’s --node-rank and --master-addr. The exec_all helper runs
torchrun on every node:
version: v4
# Multi-node PyTorch/NCCL training using the `generic` multinode implementation.
#
# NCCL is not a launcher, so there is no dedicated NCCL implementation. The
# `generic` implementation hands torchrun the per-node rank ($RANK) and the head
# node address ($MULTINODE_NODE_IP) it needs; PyTorch then uses NCCL underneath.
jobs:
torch-train:
image:
uri: docker://nvcr.io/nvidia/pytorch:25.09-py3
env:
- NCCL_DEBUG=INFO
# For InfiniBand/RoCE, pin the fabric so NCCL does not fall back to TCP
# (find the device name with `ibv_devinfo` on a node):
# - NCCL_IB_HCA=mlx5
script: |
#!/bin/sh
. /multinode/generic-helper
# Number of nodes = entries in the hostlist.
NNODES=$(echo "$MULTINODE_HOSTLIST_NOSLOTS" | awk -F, '{print NF}')
# Run torchrun on every node. Each node takes its --node-rank from $RANK;
# all rendezvous at the head node ($1 == MULTINODE_NODE_IP). Replace
# /workspace/train.py with your training script.
exec_all - "$MULTINODE_NODE_IP" "$NNODES" <<'EOF'
torchrun \
--nnodes="$2" \
--node-rank="$RANK" \
--nproc-per-node=2 \
--master-addr="$1" \
--master-port=29500 \
/workspace/train.py
EOF
resource:
cpu:
cores: 2
affinity: NUMA
memory:
size: 32GB
devices:
nvidia.com/gpu: 2 # per node; --nproc-per-node must match
multinode:
nodes: 2
implementation: generic
Each node launches --nproc-per-node processes (one per GPU) and rendezvouses at
the head node; PyTorch coordinates the rest through NCCL. See
Distributed Workflows with Generic Multinode
for the full generic environment and helper-function reference.
Because Fuzzball runs jobs as an unprivileged user, everything must be baked into the image – you cannot install packages at runtime:
- CUDA, NCCL, and your framework (e.g. PyTorch). The NGC PyTorch images bundle
all three, plus
torchrun. - For InfiniBand/RoCE, the NCCL network plugin plus
libibverbs/libfabricmatching your cluster’s fabric.
The GPU code in the image must be built for your GPU’s architecture. A CUDA
toolkit older than the GPU fails at kernel launch with Cuda failure 'invalid argument'.
On a fabric-equipped cluster, pin the device so NCCL uses RDMA instead of falling back to TCP:
env:
- NCCL_IB_HCA=mlx5
If NCCL cannot find a high-speed device it falls back to sockets, visible in the
job log as NET/IB : No device found followed by Using network Socket.
Fuzzball exposes /dev/infiniband to multinode containers automatically and
raises the locked-memory limit RDMA pinning needs; the fabric libraries must be
present in your container, and the interconnect provisioned on the nodes.
For NCCL programs that use MPI to bootstrap (for example NVIDIA’s
nccl-tests), launch them with mpirun on
the generic implementation rather than torchrun. The
generic command example
shows the pattern – your command invokes mpirun -H $MULTINODE_HOSTLIST --mca plm_rsh_agent $MULTINODE_SSH_WRAPPER -np $MULTINODE_TOTAL_SLOTS <program>, and
the program uses NCCL for the actual GPU communication.