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.
- 7 min read
- 3 reading levels
- Published
Read these first
On this page 5
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
- Your first DDP run with torchrun — the full training script, launched properly.
- DistributedSampler and set_epoch — how each worker gets its own slice of data.
- all_reduce, all_gather and broadcast — the conversation primitives underneath.
Developer — Code and libraries.
Setup
pip install torchDataParallel 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
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 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:
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()torchrun --nproc_per_node=2 ddp_hello.pyrank 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
| DataParallel | DDP | |
|---|---|---|
| processes | 1 | one per GPU |
| model copies | re-sent every step | permanent, one per process |
| per-step traffic | inputs, outputs, model | gradients only |
| Python bottleneck | one thread for all GPUs | none |
| multiple machines | no | yes |
| effort | one line | launcher + 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
- Your first DDP run with torchrun — the full training script, launched properly.
- DistributedSampler and set_epoch — how each worker gets its own slice of data.
- all_reduce, all_gather and broadcast — the conversation primitives underneath.
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
- Your first DDP run with torchrun — the full training script, launched properly.
- DistributedSampler and set_epoch — how each worker gets its own slice of data.
- all_reduce, all_gather and broadcast — the conversation primitives underneath.