DDP 深度教程
1. DDP 到底解决什么问题
DDP(DistributedDataParallel)解决的核心问题非常直接:
让多个 GPU 分别计算不同数据上的梯度,然后把这些梯度同步起来,使所有 GPU 可以像在一个更大的 batch 上训练一样更新同一个模型。
假设现在有 4 张 GPU:
GPU0 GPU1 GPU2 GPU3 │ │ │ │ │ Batch 0 │ Batch 1 │ Batch 2 │ Batch 3 ▼ ▼ ▼ ▼ Forward Forward Forward Forward │ │ │ │ ▼ ▼ ▼ ▼ Backward Backward Backward Backward │ │ │ │ ▼ ▼ ▼ ▼ g0 g1 g2 g3 │ │ │ │ └──────────────┬──────┴──────────────┬──────┴──────────────┬──────┘ │ │ │ └─────────────────────┴─────────────────────┘ │ AllReduce │ ▼ 相同的 global gradient │ ┌───────────┼───────────┐ ▼ ▼ ▼ GPU0 GPU1 GPU2 ... GPU3 │ │ │ └───────────┴───────────┘ │ Optimizer.step()因此 DDP 最重要的不是:
“怎么启动 4 个进程。”
而是:
4 个 Rank 各自计算什么?计算完以后为什么需要通信?通信以后为什么所有 Rank 可以独立更新自己的模型?
整个 DDP 可以围绕两条主线理解:
DDP │ ┌──────────┴──────────┐ │ │ 数据分发 梯度同步 │ │ ▼ ▼ 每个 Rank 拿不同数据 每个 Rank 本地计算 gradient │ │ ▼ ▼ Forward Backward │ ▼ Autograd │ ▼ Gradient Hook │ ▼ Gradient Bucket │ ▼ AllReduce │ ▼ 每个 Rank 得到相同 gradient │ ▼ Optimizer.step() │ ▼ 各 Rank 参数继续保持一致PyTorch 官方文档也明确指出,DistributedDataParallel本身不会负责切分输入数据;数据如何分布给不同 GPU,需要由用户通过DistributedSampler或其他 distributed data pipeline 完成。DDP 自身主要负责的是模型副本之间的梯度同步。
2. DDP 的基本架构
2.1 一个 Rank 对应一个训练进程
典型的单机 4 卡 DDP:
Machine 0 Process 0 ── GPU0 Process 1 ── GPU1 Process 2 ── GPU2 Process 3 ── GPU3这里有几个概念需要分清:
| 概念 | 含义 |
|---|---|
world_size | 全部参与训练的 process 数量 |
rank | 全局唯一的 process ID |
local_rank | 当前机器内部的 process / GPU ID |
例如:
2 machines × 4 GPUs Node 0: local_rank 0 → rank 0 local_rank 1 → rank 1 local_rank 2 → rank 2 local_rank 3 → rank 3 Node 1: local_rank 0 → rank 4 local_rank 1 → rank 5 local_rank 2 → rank 6 local_rank 3 → rank 7所以:
world_size = 8而:
local_rank = 2只表示:
“我是当前机器的第 2 个 GPU 对应进程。”
它不是整个训练集群中的唯一编号。
2.2 每个 Rank 都有一份完整模型
DDP 的基本模式是:
Rank 0 Rank 1 Rank 2 完整 Model 完整 Model 完整 Model │ │ │ ▼ ▼ ▼ Batch 0 Batch 1 Batch 2因此 DDP 的模型参数在正常情况下是 replicated 的,而不是 shard 的。
例如模型参数为:
θ = [10 GB]4 卡 DDP:
GPU0 → 10 GB GPU1 → 10 GB GPU2 → 10 GB GPU3 → 10 GB所以 4 卡 DDP 并不会让一个 10 GB 模型自动变成每卡 2.5 GB。
这也是 DDP 与 FSDP / ZeRO 的一个根本区别:
DDP 每张 GPU: 一份完整参数 一份自己的 gradient 一份 optimizer state因此 DDP 很适合模型本身能够放进单张 GPU 的情况。
3. 数据如何分发
3.1 DDP 不负责切数据
假设训练集:
A B C D E F G H4 个 Rank。
理想情况下:
Rank0 → A B Rank1 → C D Rank2 → E F Rank3 → G H而不是:
Rank0 → A B C D Rank1 → A B C D Rank2 → A B C D Rank3 → A B C D如果 4 张卡处理完全一样的数据,那么 DDP 只是重复计算,没有产生 data parallelism 的收益。
PyTorch 官方明确说明,DDP 不会自动 chunk / shard input,通常需要通过DistributedSampler等机制完成数据切分。
3.2 DistributedSampler 做什么
典型代码:
sampler=DistributedSampler(dataset,num_replicas=world_size,rank=rank,shuffle=True,)loader=DataLoader(dataset,batch_size=batch_size,sampler=sampler,)这里实际上建立了一条关系:
Dataset ↓ DistributedSampler ↓ 当前 Rank 对应的 sample indices ↓ DataLoader ↓ local batch例如:
global dataset 0 1 2 3 4 5 6 7 8 9 10 11 ... ↓ DistributedSampler ┌──────┼──────┬──────┐ ▼ ▼ ▼ ▼ Rank0 Rank1 Rank2 Rank3 │ │ │ │ ▼ ▼ ▼ ▼ local local local local data data data data注意一个容易混淆的地方:
DistributedSampler是负责“哪些样本属于哪个 Rank”,而DataLoader负责“如何把这些 sample 组织成 batch,并完成读取、worker、prefetch 等数据管道工作”。
3.3set_epoch()为什么重要
训练中通常会:
forepochinrange(num_epochs):sampler.set_epoch(epoch)forbatchinloader:...原因不是为了告诉 sampler “现在是第几轮”这么简单,而是:
让不同 epoch 使用不同的 shuffle 顺序,同时保证所有 Rank 对同一个 epoch 使用兼容的随机划分。
如果不调用set_epoch(),DistributedSampler 的 shuffle 顺序可能在不同 epoch 间保持不变,导致每轮数据排列重复。PyTorch 官方文档也要求在每个 epoch 开始时调用set_epoch()以让 shuffle 在 epoch 间变化。
因此它解决的是:
Epoch 0 shuffle(seed=0) ↓ Rank0 / Rank1 / Rank2 / Rank3 Epoch 1 shuffle(seed=1) ↓ Rank0 / Rank1 / Rank2 / Rank3而不是每个 Rank 随便各自 shuffle。
4. Local Batch、Global Batch 与 Global Gradient
这部分是理解 DDP 的数学核心。
假设:
world_size = N local batch size = B那么通常:
B_{\mathrm{global}} = N \times B例如:
4 GPUs 每 GPU batch = 32那么 effective global batch:
4 × 32 = 128注意这里的 “global batch” 指的是:
一次同步更新中,所有 Rank 总共参与计算的数据规模。
4.1 每个 Rank 先独立计算自己的 loss
设第i个 Rank 的 local batch 为:
D_i那么它计算:
L_i = L(D_i,\theta)然后计算本地 gradient:
g_i = \nabla_\theta L_i于是 4 个 Rank:
Rank0 → g0 Rank1 → g1 Rank2 → g2 Rank3 → g3由于它们看到的数据不同:
D0 ≠ D1 ≠ D2 ≠ D3所以一般:
g0 ≠ g1 ≠ g2 ≠ g34.2 那么真正想要的 global gradient 是什么?
如果 global batch 是所有 local batch 的并集,那么希望得到:
g = \frac{1}{N} \sum_{i=1}^{N}g_i即:
g = average(g0, g1, g2, g3)因此 DDP 真正需要解决的问题就是:
如何让所有 Rank 从自己的 local gradient 出发,最终得到相同的 global gradient?
这就是 Gradient Synchronization。
5. 为什么 gradient 必须同步
假设没有梯度同步:
Rank0 → g0 → optimizer → θ0 Rank1 → g1 → optimizer → θ1 Rank2 → g2 → optimizer → θ2 Rank3 → g3 → optimizer → θ3第一次更新之后:
θ0 ≠ θ1 ≠ θ2 ≠ θ3那么下一轮:
Rank0 用 θ0 Rank1 用 θ1 Rank2 用 θ2 Rank3 用 θ3这时它们实际上已经不是在训练同一个模型了。
DDP 希望得到的是:
θ0 = θ1 = θ2 = θ3然后下一轮继续:
不同数据 ↓ 不同 local gradient ↓ 同步 ↓ 相同 global gradient ↓ 各 Rank 对自己的同一份 θ 做相同更新 ↓ 仍然保持参数一致所以 DDP 的关键逻辑可以写成:
不同 data ↓ 不同 local gradient ↓ Gradient Synchronization ↓ 相同 gradient ↓ 相同 optimizer update ↓ 相同 parameters6. 为什么同步 Gradient,而不是直接同步 Parameter
这是面试非常容易被问到的问题。
看起来似乎可以:
Rank0: Forward Backward Optimizer.step() ↓ 得到 θ0 然后: AllReduce(θ0, θ1, θ2, θ3)但这不是 DDP 的核心设计。
原因在于:
参数同步发生在 optimizer update 之后,而 gradient synchronization 是多个 local training computation 汇聚成一次 global optimization step 的自然位置。
更重要的是,如果先各自 update:
θ0' = optimizer(θ, g0) θ1' = optimizer(θ, g1) θ2' = optimizer(θ, g2) θ3' = optimizer(θ, g3)然后再平均参数:
average(θ0', θ1', θ2', θ3')对于简单 SGD 的某些形式,看起来可能近似等价,但对于:
- momentum
- Adam
- AdamW
- weight decay
- adaptive optimizer state
参数平均并不等价于:
先得到 global gradient ↓ 再让每个 Rank 执行 optimizer.step()DDP 采用 gradient synchronization 的好处是:
每个 Rank: 本地计算 gradient ↓ 同步 global gradient ↓ 各自使用完全相同的 optimizer state ↓ optimizer.step()只要初始参数、optimizer state 和同步后的 gradient 保持一致,那么每个 Rank 的更新就是一致的。
因此:
DDP 的同步边界: Gradient ↑ Forward / Backward 是每个 Rank 本地完成 ↓ Optimizer 是每个 Rank 本地执行这使 DDP 不需要一个“中心参数服务器”。
7. DDP 的一次完整训练 Step
把前面的知识连接起来:
Dataset ↓ DistributedSampler ↓ Rank-local Batch ↓ Forward ↓ Loss ↓ Backward ↓ Autograd ↓ Parameter Gradient Ready ↓ DDP Hook ↓ Gradient Bucket ↓ Bucket Ready ↓ AllReduce ↓ Synchronized Gradient ↓ Optimizer.step() ↓ Parameters remain synchronized这就是整篇教程最重要的一条链。
下面开始进入真正的 DDP 底层机制。
8. Autograd 在 DDP 中扮演什么角色
8.1 Backward 不是“一次性计算出所有 gradient”
很多人第一次学习 DDP 时,会形成这样一个模型:
loss.backward() ↓ 所有参数 gradient 一次性全部算完 ↓ DDP 开始 AllReduce这会错过 DDP 最关键的性能设计。
实际上 backward 本身是一个按照 autograd graph 反向执行的过程:
Loss ↓ Layer N ↓ Layer N-1 ↓ Layer N-2 ↓ ... ↓ Layer 1随着 backward 向前推进,不同 parameter 的 gradient 会陆续 ready。
可以抽象成:
Backward │ ├── grad(param A) ready │ ├── grad(param B) ready │ ├── grad(param C) ready │ ├── grad(param D) ready │ └── ...DDP 正是利用了这个时间结构。
9. Autograd Hook:DDP 如何知道某个 Gradient Ready 了
DDP 不可能通过简单的:
loss.backward()# backward 完成以后allreduce(all_grads)实现最佳性能。
因为这样会浪费很多时间。
DDP 的核心思路是:
监听参数 gradient 的 ready 状态,当某个参数的 gradient 在 backward 中准备好以后,通知 DDP reducer。
可以理解为:
Autograd │ ├── Parameter A gradient ready │ └──────→ DDP Hook │ ▼ A 对应的 bucket 状态更新当某一个 bucket 中的所有参数 gradient 都 ready:
Parameter A → ready Parameter B → ready Parameter C → ready Parameter D → ready ↓ Bucket 0 ready ↓ AllReduce(Bucket 0)与此同时:
Backward 继续计算后面的 gradient于是通信和计算就可以重叠。
这就是:
Autograd → Hook → Bucket → AllReduce → Overlap
PyTorch DDP 官方文档明确说明,DDP 会把参数组织成多个 bucket,使 bucket 的 gradient reduction 能够和 backward computation 尽可能重叠;当一个 bucket 中对应的 gradient 都可用时,就可以发起该 bucket 的 reduction。
10. 为什么不能对每个 Parameter 单独 AllReduce
假设模型有:
10000 个 parameter tensors最朴素的方式: