Fuzzball Documentation
Toggle Dark/Light/Auto mode Toggle Dark/Light/Auto mode Toggle Dark/Light/Auto mode Back to homepage

Distributed GPU Training with NCCL

Overview

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.

Single node, multiple GPUs

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.

Multiple nodes

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.

What your container must provide

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/libfabric matching 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'.

InfiniBand / RDMA

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.

NCCL without PyTorch

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.