Multi-GPU and Distributed Training

DataParallel vs DistributedDataParallel

DataParallel is one process bossing several GPUs from a single thread; DistributedDataParallel runs one full worker per GPU — which is why DDP is faster and DP is effectively retired.

On this page 5
  1. Why two tools exist
  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.

Both tools train one model on several GPUs at once; the old one uses a single boss shouting instructions, the new one hires a full worker per GPU.

Picture two ways to run four cooking stations. Way one: a single head cook stands in the middle and chops nothing himself. He runs to each station carrying ingredients, instructions and finished dishes back and forth. Way two: four full cooks, one per station, each with the same recipe book, pausing once per dish to agree on seasoning.

The first way has one exhausted man as the bottleneck. The second way scales.

Why two tools exist

DataParallel (the boss) came first because it is one line of code and needs no special launching. Each training step, the boss copies the model to every GPU, splits the batch, gathers all results back to himself, and computes corrections centrally. All that copying happens every step, and one Python process coordinates everything single-handedly.

DistributedDataParallel (DDP, the four cooks) starts one complete Python process per GPU. Each process keeps its own permanent model copy and its own slice of data. Once per step, the processes share their correction signals with each other — that is the only conversation. No boss, no per-step model copying.

How it works

DataParallel (one process):
  boss GPU: [model] --copy every step--> GPU1, GPU2, GPU3
            gather all outputs back, compute, scatter again

DDP (one process per GPU):
  GPU0 [model] <----\
  GPU1 [model] <----+---- share corrections once per step
  GPU2 [model] <----+     (models never travel)
  GPU3 [model] <----/

A real example you have seen

Group projects. One "leader does everything, others wait" group finishes late — that is DataParallel. A group where everyone owns a chapter and meets briefly to align — that is DDP. Same people, different structure, different speed.

Remember this

  • Both split each batch across GPUs; the difference is who coordinates.
  • DataParallel: one process, model copied every step, one GPU overloaded.
  • DDP: one process per GPU, only corrections travel — use DDP, always.

What to learn next

Developer — Code and libraries.

Setup

bash
pip install torch

DataParallel you should recognise, not adopt; PyTorch's own documentation says use DDP instead. Captured with torch 2.5.1.

DataParallel: the one-liner and its ceiling

dp_demo.py
import torch
import torch.nn as nn

model = nn.Linear(8, 2)
dp = nn.DataParallel(model)               # one line, no launcher needed

x = torch.randn(6, 8)
print("output shape:", tuple(dp(x).shape))
print("state_dict keys:", list(dp.state_dict()))
print("GPUs it will split across:", dp.device_ids)
Output
output shape: (6, 2)
state_dict keys: ['module.weight', 'module.bias']
GPUs it will split across: [0]

It wraps anything and runs anywhere (on this single-GPU machine it has one device and splits nothing). Note the module. prefix stamped onto every checkpoint key — a loading headache both wrappers share. With several GPUs, each forward scatters the batch, replicates the model, and gathers outputs — all through one process whose Python thread becomes the choke point.

DDP: one worker per GPU

DDP code looks nearly identical, but how you run it changes — a launcher starts several copies of your script:

ddp_hello.py
import torch
import torch.distributed as dist

dist.init_process_group(backend="gloo")   # gloo works on CPU; nccl needs GPUs
rank = dist.get_rank()
world = dist.get_world_size()

t = torch.tensor([float(rank + 1)])
dist.all_reduce(t, op=dist.ReduceOp.SUM)
print(f"rank {rank} of {world}: after all_reduce ->", t.item())

dist.destroy_process_group()
bash
torchrun --nproc_per_node=2 ddp_hello.py
Output
rank 0 of 2: after all_reduce -> 3.0
rank 1 of 2: after all_reduce -> 3.0

(Two processes print concurrently, so line order varies run to run.) Each process gets a rank — its worker number — and a world size — the team headcount. all_reduce is the "share and agree" step: both workers contributed a number, both ended holding the sum. DDP's gradient sharing is this operation, run on your gradients automatically. No GPU required to learn it: the gloo backend runs the whole dance on CPU, which is how this output was captured.

The full training script, with the model wrapped and data split, is the next lesson.

Why DDP wins, concretely

DataParallelDDP
processes1one per GPU
model copiesre-sent every steppermanent, one per process
per-step trafficinputs, outputs, modelgradients only
Python bottleneckone thread for all GPUsnone
multiple machinesnoyes
effortone linelauncher + a few lines

The one-thread point deserves respect: Python largely runs one thread at a time, so DataParallel's "parallel" GPUs still queue behind a single interpreter for coordination.

Common mistakes

Reaching for DataParallel because it is less code. The line count difference is ten minutes of work; the speed difference on 4 GPUs is often 20–40%, plus DataParallel caps out at one machine.

Testing multi-GPU logic only on multi-GPU machines. The gloo backend runs DDP semantics on plain CPUs — two terminal processes on a laptop. Every lesson in this section exploits that.

Forgetting the checkpoint prefix. Both wrappers save keys as module.weight. Save model.module.state_dict() to keep checkpoints clean — details in saving from one rank.

Expecting a 4x speedup from 4 GPUs. Communication takes time. Well-tuned DDP reaches 3.5x-ish on 4 local GPUs for meaty models; tiny models communicate more than they compute.

Try it yourself

Run ddp_hello.py with --nproc_per_node=4 and predict the printed sum before it runs. Then change SUM to dist.ReduceOp.MAX and predict again.

What to learn next

Researcher — Mathematics and papers.

The two designs, precisely

DP: single-process, multi-thread. Per iteration: replicate parameters to each device, scatter input along batch dim, parallel forward via threads, gather outputs to device 0, compute loss there, scatter gradients back through backward, reduce parameter gradients onto the primary replica, update, repeat. Costs: model broadcast $O(P)$ per step, output gather at device 0, GIL-serialised kernel launching, and device-0 memory/computation asymmetry.

DDP: multi-process SPMD. Parameters broadcast once at construction. Each backward computes local gradients; DDP averages them across ranks with all-reduce so every replica applies identical updates and replicas never diverge. Communication per step is exactly one gradient's worth, $O(P)$, but overlapped: parameters are grouped into buckets (default 25 MB), and each bucket's all-reduce launches as soon as its gradients are ready during backward — communication hides behind remaining computation. The Reducer uses autograd hooks per parameter; bucket order approximates reverse registration order, and find_unused_parameters=True handles graphs where some parameters get no gradient (at the cost of a graph walk per step — see the sync lesson).

Ring all-reduce moves $2P \frac{N-1}{N}$ bytes per rank regardless of $N$ (bandwidth-optimal); NCCL supplies ring and tree variants over NVLink/PCIe/network. Scaling efficiency is then governed by the ratio of backward compute time to gradient transfer time — large conv/transformer blocks overlap almost fully; tiny MLPs do not.

Historical note and the modern menu

DP descends from the Krizhevsky-era single-node pattern; DDP's overlapped-bucket design is described in Li et al. (2020). Beyond DDP, the same process-per-device substrate carries FSDP (shard parameters too), tensor/pipeline parallelism (model parallel strategies), and their compositions — DDP remains the baseline all of them are measured against.

References

  • Li et al. (2020), PyTorch Distributed: Experiences on Accelerating Data Parallel Training, VLDB — the DDP design paper: buckets, overlap, and scaling measurements.
  • Sergeev and Del Balso (2018), Horovod: fast and easy distributed deep learning in TensorFlow — popularised ring all-reduce for DL.
  • Patarasuk and Yuan (2009), Bandwidth optimal all-reduce algorithms for clusters of workstations — the ring result.

What to learn next