跳到主要内容

数据截至 (上游 commit c187ef3271d5)

05 · 分布式训练:DDP 与梯度桶

这一章讲什么: 从单机到多卡到底多了什么。读完你会理解:为什么 DDP 不需要你改模型代码、梯度同步为什么几乎不占用训练时间、以及 find_unused_parameters 那个报错是从哪来的。


1. 它要解决的小问题

数据并行的原理一句话:每张卡拿不同数据、跑同一个模型,反向后把梯度跨卡求平均,再各自更新——只要梯度一致,各卡模型就永远一致。

朴素实现的两个性能问题:

  • 等所有梯度算完才统一同步 → 通信时间纯浪费;
  • 逐参数发起通信 → 几千个小 allreduce 的延迟远大于带宽。

DDP(DistributedDataParalleltorch/nn/parallel/distributed.py:466)用两个机制同时解掉:梯度桶(聚小成大)+ autograd 钩子(边算边传)。


2. 地基:ProcessGroup 通信原语

一切集合通信的抽象是 ProcessGrouptorch/csrc/distributed/c10d/ProcessGroup.hpp:78),虚接口 allreduceProcessGroup.hpp:256)等由各后端(NCCL/Gloo/MPI)实现。Python 侧入口是 init_process_grouptorch/distributed/distributed_c10d.py:2374)和 all_reducetorch/distributed/distributed_c10d.py:3956)。

DDP 与 FSDP 都建在这层之上——通信怎么发是 c10d 的事,什么时候发是 DDP 的事。


3. DDP 的一生:初始化、前向、反向

3.1 初始化:先对齐,再挂钩子

构造时做两件大事:

  1. 参数对齐:rank 0 的参数广播给所有 rank(_broadcast_coalescedtorch/nn/parallel/distributed.py:1224)——保证大家从同一个初始权重出发;
  2. 建 Reducer_ddp_init_helpertorch/nn/parallel/distributed.py:1375)里创建 C++ 的 dist.Reducertorch/nn/parallel/distributed.py:1438),按参数逆序分桶(initialize_bucketstorch/csrc/distributed/c10d/reducer.hpp:72)。

3.2 钩子在哪:挂在梯度累加器上

关键机制在 C++ Reducer 构造里(torch/csrc/distributed/c10d/reducer.cpp:185-202):对每个参数,拿到它的 grad_accumulator(第 3 章的 AccumulateGrad 节点),注册一个 post hook——

// torch/csrc/distributed/c10d/reducer.cpp:191-202(节选)
hooks_.emplace_back(
grad_accumulator->add_post_hook(std::make_unique<...LambdaPostHook>(
[this, variable_index](const variable_list& outputs, ...) {
...
this->autograd_hook(variable_index); // 梯度一就绪就通知 Reducer
return outputs;
}, ...)),
grad_accumulator);

也就是说:DDP 完全不改动你的前向计算,它只是买通了 autograd 引擎的门房——每个参数的梯度一进 .grad,门房就打电话来。

3.3 反向时:桶装满就发车

Reducer::autograd_hooktorch/csrc/distributed/c10d/reducer.cpp:668)注释写得很清楚:它在「某个模型参数的梯度被累加完之后」被调,且只从 autograd 线程调用。它把该梯度拷入对应桶;桶内所有梯度到齐,就对这个桶发起一次 allreduce(all_reduce_buckettorch/csrc/distributed/c10d/reducer.cpp:977)。全部桶发完后 finalize_backwardtorch/csrc/distributed/c10d/reducer.cpp:658)做收尾(等待、写回、unused 参数检查)。

时间轴上的效果(重叠的来源):

反向计算: [层N grad] [层N-1 grad] ... [层1 grad]
│ │ │
桶就绪: 桶2装满 ──► allreduce桶2 桶1装满 ──► allreduce桶1
(通信与后续反向计算同时进行)

因为桶按参数逆序分(输出侧参数先进桶),最早算好的梯度最早发车。

3.4 前向也要配合

_pre_forwardtorch/nn/parallel/distributed.py:1745)里 reducer.prepare_for_forward():1759)重置本轮状态;_post_forwardprepare_for_backward(list(_find_tensors(output))):1838)告诉 Reducer 本次输出涉及哪些张量——用于 unused 参数检测。


4. 原理演示

# 示意,非源码:DDP 的用法与普通模型无异
model = MyModel().to(local_rank)
ddp = DistributedDataParallel(model, device_ids=[local_rank])

loss = ddp(batch).loss # 前向:与单机完全相同
loss.backward() # 反向:梯度在后台已被跨卡平均
optimizer.step() # 各卡用「一致的平均梯度」各自更新

注意心智模型:DDP 同步的是梯度,不是参数;参数一致性是「初始一致 + 每步更新一致」的推论。


5. 坑与边界

  • unused 参数是头号地雷。 某次前向没用到的参数,其钩子不会触发,桶永远装不满;反向结束时 Reducer 逐参数检查梯度是否到齐,没收到梯度的参数直接报错("Parameter indices which did not receive grad for rank …",torch/csrc/distributed/c10d/reducer.cpp:2166)。对策:find_unused_parameters=True(动态图模式逐轮检测),或确认后开 static_graph=True(第一轮就记录未用参数集合,torch/csrc/distributed/c10d/reducer.cpp:625-644)。
  • 梯度累加(gradient accumulation)要用 no_sync() 上下文包住中间步,否则每个 micro-step 都白做一次 allreduce。
  • DDP 只解决数据并行。 模型大到单卡装不下时要用 FSDP(把参数本身分片,torch/distributed/fsdp/)或张量并行;那些是另一套取舍,本章不展开。
  • 进程组初始化依赖网络环境(store、rank、world_size),NDEBUG/超时问题大多出在 c10d 层而不是 DDP 层。
  • 看不出来:NCCL 后端内部的流管理与传输调优——在 third_party 与 c10d 后端实现里,超出本章主线。

6. 本章代码地图

主题文件路径符号名
通信抽象torch/csrc/distributed/c10d/ProcessGroup.hppProcessGroup::allreduce
Python 通信 APItorch/distributed/distributed_c10d.pyinit_process_groupall_reduce
DDP 主体torch/nn/parallel/distributed.pyDistributedDataParallel_ddp_init_helper_pre_forward_post_forward
梯度归约器torch/csrc/distributed/c10d/reducer.cppReducer::autograd_hookall_reduce_bucketfinalize_backward
分桶torch/csrc/distributed/c10d/reducer.hppReducer::initialize_buckets
初始参数广播torch/nn/parallel/distributed.py_broadcast_coalesced(:1224)