Multi-GPU and Distributed Training

Training across several machines

Your training script does not change when you add a second machine — what changes is the launch command, the network between the nodes, and which of RANK and LOCAL_RANK you were quietly getting away with confusing.

On this page 5
  1. Why the network becomes the whole story
  2. How it works
  3. A real example you have seen
  4. Remember this
  5. What to learn next

One lesson, three depths. Pick the one that fits you today — you can switch any time.

Beginner — No maths. Plain English.

Going from one machine to several changes the launch command and the network, not your training code.

Think of a cricket team practising in one net. Add a second ground across town and the drills stay the same. What changes is that everyone needs to know which ground to go to, and the two grounds need a phone line between them.

Your script is the drill. torchrun on each machine is the coach. The phone line is your network, and it is the part that decides whether this was a good idea.

Why the network becomes the whole story

Inside one machine, GPUs talk over cables built for the purpose, moving hundreds of gigabytes a second. Between two machines they talk over the office network, which may move one gigabyte a second.

That is a hundred-fold slowdown on every conversation. The workers still have to agree on their corrections after every step, and now that agreement crosses a slow wire.

This is why a second machine can make training slower. Small models, which talk often and compute little, suffer most.

How it works

machine A                          machine B
  torchrun --nnodes=2 --node_rank=0   torchrun --nnodes=2 --node_rank=1
    RANK=0 LOCAL_RANK=0                 RANK=2 LOCAL_RANK=0
    RANK=1 LOCAL_RANK=1                 RANK=3 LOCAL_RANK=1
           \                                    /
            \______ the network between them __/
              every step, all four must agree

RANK is your number in the whole team. LOCAL_RANK is which GPU you sit on, inside your own machine. On one machine they are the same, which is why the difference stays hidden until the day it bites.

A real example you have seen

A wedding with the ceremony at one hall and the reception at another. Everything works, but the two halls must coordinate by phone, and every phone call slows the day down. One hall would have been simpler if it were big enough.

Remember this

  • The training script is unchanged; the launch command and the network change.
  • RANK is your number in the whole job; LOCAL_RANK is your GPU on this machine.
  • Between machines the network is slow, so a second machine can make things worse.

What to learn next

Developer — Code and libraries.

Setup

bash
pip install torch

The output below was captured on one machine, with the environment variables forced to describe a two-node, two-process-per-node layout. That is enough to show how the numbering works. Every timing claim about real networks in this lesson comes from published measurements, not from a run performed here.

The variables torchrun gives you

topology.py
import os
import torch
import torch.distributed as dist

dist.init_process_group(backend="gloo")      # use nccl on GPU nodes

rank = int(os.environ["RANK"])               # my number across ALL machines
local_rank = int(os.environ["LOCAL_RANK"])   # my GPU index on THIS machine
world = int(os.environ["WORLD_SIZE"])        # total processes everywhere
local_world = int(os.environ["LOCAL_WORLD_SIZE"])   # processes on this machine
node = rank // local_world                   # which machine I am on

print(f"node={node}  rank={rank}/{world}  local_rank={local_rank}/{local_world}")

t = torch.tensor([float(rank)])
dist.all_reduce(t, op=dist.ReduceOp.SUM)     # crosses machines, identical call
if rank == 0:
    print(f"sum of all ranks, every machine included: {t.item():.0f}")

dist.destroy_process_group()

On two real machines, you would launch it like this:

bash
# on machine A (the one whose address everyone rendezvouses at)
torchrun --nnodes=2 --node_rank=0 --nproc_per_node=2 \
         --master_addr=10.0.0.11 --master_port=29500 topology.py

# on machine B, at the same time
torchrun --nnodes=2 --node_rank=1 --nproc_per_node=2 \
         --master_addr=10.0.0.11 --master_port=29500 topology.py
Output
node=0  rank=0/4  local_rank=0/2
node=0  rank=1/4  local_rank=1/2
node=1  rank=2/4  local_rank=0/2
node=1  rank=3/4  local_rank=1/2
sum of all ranks, every machine included: 6

Read the two right-hand columns together. Rank 2 has local_rank=0 — it is the first GPU on the second machine. Send a tensor to cuda:2 there and you address a device that may not exist; send it to cuda:0 on the wrong machine and you overload one card.

0 + 1 + 2 + 3 = 6, so the all_reduce genuinely covered all four processes. On real nodes, add socket.gethostname() to the print line to confirm the placement is what you asked for.

The rule that prevents most multi-node bugs

python
local_rank = int(os.environ["LOCAL_RANK"])
torch.cuda.set_device(local_rank)      # LOCAL_RANK selects the device
model = model.to(local_rank)
ddp = DDP(model, device_ids=[local_rank])

if int(os.environ["RANK"]) == 0:       # RANK decides who writes and logs
    torch.save(...)

LOCAL_RANK chooses hardware. RANK chooses responsibility. Every multi-node placement bug is a violation of that one sentence.

Elastic rendezvous, which you want instead

Hard-coding --master_addr means the job dies if that machine reboots. The c10d rendezvous backend is the better default:

bash
# identical command on every node -- no node_rank, no master_addr
torchrun --nnodes=2 --nproc_per_node=2 \
         --rdzv-backend=c10d --rdzv-endpoint=10.0.0.11:29500 \
         --rdzv-id=my-run-42 topology.py

Same command everywhere is a real operational win: one line in your job template rather than a per-host variant. With --nnodes=2:4 the job also starts on two nodes and absorbs more as they appear, restarting the worker group each time — which only helps if your script resumes from a checkpoint.

Under Slurm

Most shared clusters use Slurm, and the pattern is one torchrun per node:

train.sbatch
#!/bin/bash
#SBATCH --nodes=2
#SBATCH --gpus-per-node=8
#SBATCH --ntasks-per-node=1          # ONE task per node; torchrun spawns the rest

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

srun torchrun --nnodes=$SLURM_NNODES --nproc_per_node=8 \
     --node_rank=$SLURM_NODEID --master_addr=$MASTER_ADDR \
     --master_port=$MASTER_PORT train.py

--ntasks-per-node=1 is the line people get wrong. Set it to 8 and Slurm starts 8 torchruns per node, each starting 8 workers — 64 processes fighting over 8 GPUs.

When a second machine is worth it

Be honest with yourself before booking the second node.

situationsecond node likely to help?
model does not fit on one nodeyes — this is the real reason
large model, InfiniBand between nodesyes
small model, ordinary Ethernetoften no
dataloader is the bottleneckno — fix the loader first
one node with idle GPUsno — fill that node first

The tools that make multi-node bearable are all about talking less: no_sync accumulation, gradient compression through DDP communication hooks, and hybrid sharding that keeps the expensive collectives inside a node.

Common mistakes

Using RANK where LOCAL_RANK belongs. On one node they are equal and the code is fine. On two nodes, rank 2 tries to use cuda:2 and everything falls apart. The most common multi-node bug there is.

A firewall between the nodes. init_process_group hangs until the timeout with no useful message. Test with nc -zv <master_addr> 29500 from the other node before you blame PyTorch.

Different environments on each machine. Different PyTorch versions, different CUDA, different code. Symptoms range from an immediate NCCL version error to silently wrong gradients. Ship a container.

Leaving NCCL to guess the interface. On machines with several network cards, NCCL may pick the slow one. export NCCL_SOCKET_IFNAME=eth0 and export NCCL_DEBUG=INFO at the first sign of trouble; the debug output names the transport it selected.

Assuming your shared filesystem is fast. Every rank opening the same dataset over NFS creates a thundering herd. Stage data to local disk in a rank-0 step guarded by a barrier().

Try it yourself

Run topology.py on one machine with --nproc_per_node=4 under plain torchrun. Then set LOCAL_WORLD_SIZE=2 by hand and rerun to see the two-node numbering appear, exactly as the output above was captured.

What to learn next

Researcher — Mathematics and papers.

The rendezvous

torchrun runs one elastic agent per node. With --rdzv-backend=c10d, agents meet at a TCPStore hosted by the endpoint host, agree a round, and the store assigns each agent a GROUP_RANK. Global ranks are then handed out contiguously by group: node $k$ owns ranks $[k \cdot L, (k+1) \cdot L)$ for LOCAL_WORLD_SIZE $= L$, which is exactly the block structure the output above shows and the reason node = rank // local_world is a valid derivation.

Restart semantics are group-level. On any worker failure the agent kills its local workers, re-enters the rendezvous, and the whole world restarts from TORCHELASTIC_RESTART_COUNT + 1 with possibly different ranks. Nothing is preserved in memory, which makes checkpoint cadence the real determinant of usable throughput on unreliable clusters. The --max-restarts budget should be set against the expected failure rate, and the checkpoint interval against the restart cost — Meta's published guidance for large runs is checkpointing often enough that a restart costs a small single-digit percentage of progress.

What the interconnect actually costs

Ring all-reduce moves $2S\frac{N-1}{N}$ bytes per rank for a payload of $S$ bytes. The wall time is set by the slowest link in the ring, so a topology that mixes NVLink inside a node with Ethernet between nodes runs at Ethernet speed unless the algorithm is hierarchical. NCCL detects this and uses a two-level scheme — reduce-scatter within the node, all-reduce across nodes on one rank per node, all-gather back within the node — which reduces inter-node traffic by the intra-node GPU count.

Orders of magnitude worth carrying in your head: NVLink between GPUs in a node is in the hundreds of GB/s; PCIe generation 4 x16 is about 32 GB/s bidirectional; InfiniBand HDR is 25 GB/s per port; 10 GbE is about 1.2 GB/s. The gap between the first and last is roughly two orders of magnitude, which is the entire reason no_sync accumulation, gradient compression and hybrid sharding exist.

Scaling efficiency is governed by the ratio of backward-pass compute time to gradient transfer time. Since compute scales with parameters times batch and transfer scales with parameters alone, larger per-rank batches improve multi-node efficiency — an argument that pushes against the large-batch generalisation limits from the other side.

Failure modes specific to many nodes

At scale, the dominant hazards are stragglers and silent data corruption rather than crashes. A single slow rank sets the pace of every collective, so per-rank step-time histograms are the standard monitoring artefact. NCCL's watchdog, controlled by TORCH_NCCL_ASYNC_ERROR_HANDLING and the collective timeout, converts a hung collective into a process abort with a rank identified — without it, a job can sit at 100% GPU utilisation making no progress. Flight recorder dumps (TORCH_NCCL_TRACE_BUFFER_SIZE) capture the last collectives issued per rank and are the tool of choice for post-mortem on a hang; see debugging distributed hangs.

References

  • PyTorch documentation, Torch Distributed Elastic and torchrun — rendezvous, restart and the environment contract.
  • NVIDIA, NCCL Developer Guide — hierarchical algorithms, topology detection, NCCL_SOCKET_IFNAME and NCCL_DEBUG.
  • Narayanan et al. (2021), Efficient Large-Scale Language Model Training on GPU Clusters Using Megatron-LM, SC21 — measured scaling across thousands of GPUs and the intra-node/inter-node placement rules.
  • Zhao et al. (2023), PyTorch FSDP, VLDB 16(12) — hybrid sharding as a response to slow inter-node links.

What to learn next