六:分布式训练
DataParallel和DistributedDataParallel的根本区别是什么?¶
DataParallel(DP)和DistributedDataParallel(DDP)是PyTorch中两种实现数据并行的方式,但它们在设计哲学和实现机制上有本质的不同。
根本区别:DP采用单进程多线程架构,在一台机器的一个Python进程中通过多线程管理多张GPU;而DDP采用多进程架构,为每一张GPU启动一个独立的Python进程,进程间通过NCCL/GLOO等集合通信库进行消息传递。
这一区别决定了它们在性能、灵活性、扩展性上的巨大差异:
| 特性 | DataParallel (DP) | DistributedDataParallel (DDP) |
|---|---|---|
| 进程模型 | 单进程,主线程控制所有GPU | 多进程,每张GPU一个独立进程 |
| 通信方式 | Python多线程 + 共享内存(性能瓶颈在Python全局解释器锁GIL) | 高效的C++后端集合通信(NCCL),无GIL影响 |
| 负载均衡 | 主GPU需要收集所有梯度并进行广播,成为通信和计算瓶颈 | 所有GPU地位平等,通过All-Reduce高效同步 |
| 模型复制 | 前向传播时,每张GPU动态复制主模型;反向传播时,梯度汇聚到主GPU再更新 | 每个进程在初始化时创建一份完整的模型副本 |
| 扩展性 | 仅限单机,且随着GPU数量增加性能下降严重 | 可扩展到多机多卡,线性加速比高 |
| 性能开销 | 主GPU显存和计算压力大,速度慢 | 通信与计算可重叠,速度接近线性加速 |
深度解析:
DP模式下,每进行一次前向传播,主GPU都会将整个模型的参数通过Python线程广播(scatter)给其他GPU,处理完各自的batch后,其他GPU将梯度传回主GPU(gather),由主GPU统一进行参数更新,再将更新后的模型广播出去。整个过程存在大量的Python线程间数据搬运和同步,受GIL限制,且主GPU承担了过多的通信和存储压力。
DDP则彻底抛弃了Python多线程。它通过外部进程启动器(如torchrun)为每张GPU创建一个独立的Python进程,每个进程拥有独立的解释器和GIL。进程间利用NCCL(GPU之间)或GLOO(CPU之间)进行高效的C++实现集合通信,梯度同步通过All-Reduce直接完成,无需经过某个“主”GPU。这使得DDP在多卡和多机上的扩展性极佳。
为什么DDP比DP效率更高?¶
DDP比DP效率更高的原因主要来自以下几个方面:
- 摆脱了Python GIL(全局解释器锁)
DP是单进程多线程,Python的GIL导致在任何时刻只有一个线程能执行Python字节码。当多个GPU同时计算时,主线程的Python代码(包括梯度同步、参数更新)受限于GIL,无法真正并行。DDP为每张GPU启动独立进程,每个进程有自己的Python解释器和GIL,真正实现了CPU调度的并行化。
- 通信模式更高效
DP中,主GPU必须收集所有GPU的梯度(Gather操作),然后求平均或求和,再广播更新后的参数给所有GPU。主GPU的通信带宽成为瓶颈,且通信与计算无法重叠。DDP使用All-Reduce算法,所有GPU地位平等,直接参与梯度聚合,不存在单个瓶颈点。现代NCCL的All-Reduce实现(如Ring All-Reduce)能充分利用各卡之间的双向带宽,速度远快于DP的Scatter/Gather模式。
- 模型复制与生命周期管理
DP每次前向传播时,都需要将主GPU的模型参数通过Python线程广播给其他GPU,反向传播时又要回收梯度。这种动态的参数传输引入了大量额外的CPU和GPU通信开销。DDP在训练开始时一次性在每个进程内创建相同的模型副本,整个训练过程中模型参数驻留在各自的GPU上,无需反复广播。
- 通信与计算重叠
DDP支持在反向传播计算梯度的同时,异步执行梯度All-Reduce。即计算完某一层的梯度后,不等整个网络反向传播结束,就可以提前启动该层梯度的通信。这几乎将通信延迟完全隐藏在计算时间之后。DP则无法做到这种细粒度的重叠,因为它必须在所有梯度计算完成后再由主GPU统一处理。
- 负载均衡
DP中主GPU需要处理额外的广播和收集任务,导致其计算和显存负载高于其他GPU,造成负载不均。DDP所有进程负载完全对称,资源利用率高。
实验数据:在典型的8卡A100上训练ResNet50,DDP的吞吐量可达DP的3-5倍,且随着卡数增加,DP的性能甚至会倒退,而DDP保持线性加速。
使用DDP时,如何启动多进程训练?torchrun或torch.distributed.launch的作用。¶
DDP的多进程训练需要为每张GPU启动一个独立的Python进程。PyTorch提供了两个主要的命令行启动器:
-
torch.distributed.launch(旧版,已弃用) -
torchrun(新版,推荐使用)
torchrun的作用:
torchrun是一个进程启动器(launcher),它会自动读取环境变量(如WORLD_SIZE、RANK、LOCAL_RANK),在集群中的每个节点上启动相应数量的训练进程,并为每个进程设置正确的分布式环境变量。用户只需编写单进程的训练脚本,无需手动管理进程的创建和销毁。
使用示例:
# 单机多卡(2张GPU)
torchrun --nproc_per_node=2 train.py
# 多机多卡(每机2张GPU,共2台机器)
# 在节点0上执行:
torchrun --nproc_per_node=2 --nnodes=2 --node_rank=0 --master_addr="192.168.1.1" --master_port=12355 train.py
# 在节点1上执行:
torchrun --nproc_per_node=2 --nnodes=2 --node_rank=1 --master_addr="192.168.1.1" --master_port=12355 train.py
关键参数:
-
--nproc_per_node:每个节点上要启动的进程数(通常等于GPU数量)。 -
--nnodes:总节点数。 -
--node_rank:当前节点的编号(从0开始)。 -
--master_addr:主节点的IP地址。 -
--master_port:主节点的端口,所有节点通过它通信。
torchrun的优势:
-
自动处理环境变量的设置和传递。
-
支持容错和弹性训练(如
torchrun --rdzv_endpoint)。 -
比手动设置
MASTER_ADDR等环境变量更可靠、更不易出错。
在训练脚本中,可以通过torch.distributed.get_rank()等函数获取当前进程的身份,进行相应的初始化。
在DDP中,每个进程的模型、优化器和数据是怎样的?¶
在DDP中,每个进程拥有完全相同的模型初始副本、独立但同构的优化器、以及互不重叠的数据子集。
-
模型
-
所有进程在训练开始前,从相同的预训练权重(或随机初始化)加载模型。
-
训练过程中,由于所有进程的梯度通过All-Reduce求平均,因此它们在每一次更新后,参数保持完全一致。
-
在加载模型时,通常只在rank 0从磁盘读取权重,然后通过广播或直接让其他进程也从同一路径加载(文件系统共享)来初始化。
-
优化器
-
每个进程拥有自己的优化器实例,存储着该进程本地参数副本的优化器状态(如Adam的m和v)。
-
由于梯度全局同步,各进程的优化器状态在数学上也是同步的。
-
保存检查点时,通常每个进程都会保存自己的优化器状态,以便恢复时可以精确还原分布式状态。
-
数据
-
使用
DistributedSampler,将整个训练集均匀划分到各个进程上,每个进程看到的数据子集互不重叠。 -
在每个epoch内,所有进程总共处理了整个数据集一次(假设没有
drop_last和数据集长度可被进程数整除)。 -
不同的进程加载不同的数据,生成各自的本地梯度,通过All-Reduce获得全局梯度。
流程总结:
-
各进程独立从自己的数据子集中获取一个batch。
-
独立前向传播,计算损失。
-
独立反向传播,计算本地梯度。
-
All-Reduce同步梯度(求平均)。
-
各进程用平均梯度更新本地模型参数。
-
由于初始模型相同、梯度相同,更新后的参数仍然一致。
这种设计使得DDP在保证模型一致性的前提下,最大化了数据处理的并行度。
如何初始化分布式训练环境?init_process_group的参数含义。¶
分布式训练环境通过torch.distributed.init_process_group()进行初始化,它建立进程间通信的基础。
函数原型:
torch.distributed.init_process_group(
backend='nccl', # 通信后端
init_method='env://', # 进程发现与握手方式
world_size=4, # 总进程数
rank=0, # 当前进程的全局编号
timeout=timedelta(minutes=30) # 超时设置
)
参数详解:
backend:集合通信的后端类型。'nccl':NVIDIA的集合通信库,用于GPU间通信,是分布式训练的标准选择,提供最佳性能。'gloo':通用后端,支持CPU和GPU,但GPU性能不如NCCL,常用于CPU分布式或作为NCCL的备选。-
'mpi':基于消息传递接口,需要额外安装MPI库,较少使用。 -
init_method:指定各个进程如何互相发现并进行初始握手。 'env://':从环境变量中读取MASTER_ADDR、MASTER_PORT、WORLD_SIZE、RANK。这是配合torchrun最方便的方式。'tcp://...':手动指定一个TCP地址作为会合点,所有进程必须能连接到该地址。-
'file://...':使用共享文件系统上的一个文件作为协调点(不常用)。 -
world_size:参与分布式训练的总进程数(= GPU总数)。通常在脚本中通过torch.distributed.get_world_size()获取,而不是硬编码。 -
rank:当前进程的全局唯一编号,从0到world_size-1。同样通常从环境变量中获取。 -
timeout:集合通信操作的超时时间,默认为10分钟。对于大模型或慢速网络,可能需要调大以避免超时错误。
典型初始化代码:
import torch.distributed as dist
import os
# 当使用torchrun启动时,环境变量已设置
dist.init_process_group(backend='nccl', init_method='env://')
local_rank = int(os.environ['LOCAL_RANK'])
torch.cuda.set_device(local_rank) # 绑定当前进程到对应GPU
初始化完成后,所有进程通过一个逻辑通信组连接起来,可以进行后续的集合操作。
分布式训练中,local_rank、global_rank和world_size分别指什么?¶
这三个概念是理解分布式训练拓扑的关键。
-
world_size:参与训练的总进程数。对于单机多卡,即GPU数量;对于多机多卡,即总节点数 × 每节点GPU数。它是一个全局常量。 -
global_rank:全局唯一的进程ID,从0到world_size-1。在整个集群中,每个进程都有一个唯一的global_rank。它用于在跨节点的通信中标识进程。 -
local_rank:节点内的本地GPU编号,从0到节点内GPU数-1。它用于指定当前进程应该使用哪张物理GPU。在同一台机器上,不同进程的local_rank不同;但不同机器上的进程可以有相同的local_rank。
举例:假设有2台服务器,每台4张GPU,总共8张GPU。
-
world_size = 8 -
global_rank:节点0上的4个进程分别为0,1,2,3;节点1上的4个进程分别为4,5,6,7。 -
local_rank:在每台机器上,4个进程的local_rank都是0,1,2,3。
在代码中的使用:
import torch.distributed as dist
world_size = dist.get_world_size()
global_rank = dist.get_rank()
local_rank = int(os.environ['LOCAL_RANK'])
torch.cuda.set_device(local_rank) # 绑定本地GPU
理解这三个概念有助于正确配置数据采样、模型保存和调试分布式程序。
如何获取当前进程的rank和总进程数?¶
使用torch.distributed提供的API:
import torch.distributed as dist
# 必须已经调用过 init_process_group
global_rank = dist.get_rank() # 当前进程的全局排名
world_size = dist.get_world_size() # 总进程数
注意:这两个函数必须在init_process_group成功调用之后才能使用,否则会抛出异常。
在非分布式脚本中的兼容写法:
if torch.distributed.is_available() and torch.distributed.is_initialized():
rank = dist.get_rank()
world_size = dist.get_world_size()
else:
rank = 0
world_size = 1
这样可以保证代码在单卡和多卡环境下都能正常运行。
DistributedSampler如何保证每个进程获取不同数据?¶
DistributedSampler通过将整个数据集的索引按进程总数均匀划分,然后根据当前进程的rank分配不重叠的子集,确保每个进程处理的数据互不重复。
工作原理:
-
假设数据集有100个样本,
world_size=4。 -
DistributedSampler生成索引[0, 1, 2, ..., 99]。 -
将索引按进程数划分:rank 0 取
[0, 4, 8, ...],rank 1 取[1, 5, 9, ...],以此类推。 -
划分后的索引子集作为每个进程的采样范围。这样,整个epoch内,所有进程的数据覆盖了整个数据集,且互不重叠。
实现细节:
-
如果数据集长度不能被进程数整除,可以通过
drop_last=True丢弃最后不完整的batch,或通过pad补齐。 -
通过
shuffle=True,在每个epoch开始时,DistributedSampler会使用相同的随机种子打乱全局索引,然后分配给各进程,保证所有进程的随机性一致且数据不重叠。
使用示例:
from torch.utils.data import DistributedSampler
sampler = DistributedSampler(dataset, num_replicas=world_size, rank=rank, shuffle=True)
dataloader = DataLoader(dataset, batch_size=32, sampler=sampler)
这样,每个进程的DataLoader只会从分配给它的子集中抽取数据,无需额外设置。
为什么需要在每个epoch开始前调用sampler.set_epoch(epoch)?¶
DistributedSampler.set_epoch(epoch)的作用是为每个epoch设置不同的随机种子,以保证不同epoch中数据的打乱顺序不同,从而增强训练的随机性和泛化能力。
原理:
-
DistributedSampler在shuffle=True时,内部使用torch.Generator生成随机排列。该生成器的种子基于epoch值计算(通常是base_seed + epoch)。 -
如果不调用
set_epoch,默认epoch=0,则每个epoch的打乱顺序完全相同,数据多样性丧失,模型容易记住数据顺序而过拟合。 -
调用后,每个epoch会有一个全新的随机排列,使得模型每个epoch看到的数据顺序都不同。
使用方式:
for epoch in range(num_epochs):
train_sampler.set_epoch(epoch) # 必须在每个epoch开始前调用
for batch in train_loader:
...
注意:在分布式训练中,必须确保所有进程调用的set_epoch传入相同的epoch值,以保证全局数据打乱的一致性,从而继续维持各进程数据子集的不重叠性。
DDP中梯度是如何同步的?AllReduce发生在什么时候?¶
在DDP中,梯度同步通过All-Reduce操作完成,该操作将每个进程的本地梯度张量求和或求平均,并让所有进程得到相同的结果。
发生时机:
DDP采用自动微分的Hook机制,在反向传播过程中,每当计算完一个参数的梯度,autograd引擎就会触发为该参数注册的Hook。DDP的Hook会立即启动对该参数的异步All-Reduce。这意味着通信与计算可以重叠:当模型后面层的梯度还在计算时,前面层的梯度已经在同步了。所有同步完成后,每个进程拥有了全局平均的梯度。
同步算法的选择:
默认使用NCCL的All-Reduce,底层通常是高效的Ring All-Reduce算法。它将所有GPU组成一个逻辑环,梯度被分割成小块,在环上并行传输和累加,最后广播。这种方法通信量约为2 * (N-1) / N * 数据量,效率极高。
梯度的平均:
All-Reduce默认对梯度求和,但为了保持全局batch size的语义,DDP在backward后会自动将求和结果除以world_size,得到平均梯度。这样,无论使用多少GPU,参数更新的量级与大batch训练一致。
图示流程:
-
loss.backward()触发反向传播。 -
计算第L层梯度。
-
DDP Hook自动异步启动第L层梯度的All-Reduce。
-
继续计算第L-1层梯度,同时第L层的通信正在进行。
-
... ...
-
所有层梯度同步完成,每个进程拥有完全相同的平均梯度。
什么是“参数服务器”架构?PyTorch是否原生支持?¶
参数服务器(Parameter Server, PS) 是一种经典的分布式机器学习架构,由服务器节点和工作者节点组成。参数服务器存储全局模型参数,工作者节点负责计算本地梯度,然后将梯度推送到参数服务器,服务器更新参数后,工作者再拉取最新参数。
优缺点:
-
优点:服务器集中管理参数,易于实现异步更新,适合计算力不均衡的异构环境。
-
缺点:服务器成为通信和存储瓶颈(需要存储全部参数),且同步策略复杂,扩展性受限。
PyTorch是否原生支持?
PyTorch没有原生内置参数服务器架构。它的设计哲学是去中心化的集合通信(如DDP + NCCL),所有节点地位平等。PyTorch通过torch.distributed模块提供了基础的通信原语(send, recv, broadcast等),用户可以自己实现参数服务器逻辑,但这需要大量工作。目前主流的大规模训练已全面转向All-Reduce/ZeRO等去中心化方式,PS架构在大模型时代逐渐式微。
在DDP中,如何保存和加载模型?需要特别注意什么?¶
保存模型:
由于DDP中每个进程的模型参数是一致的,因此通常只在rank 0(主进程)上保存即可,避免多进程同时写入造成冲突。保存的是model.module.state_dict()(因为DDP将模型包裹了一层),以获取原始模型的参数。
加载模型:
通常在训练开始时,由rank 0加载权重,然后通过dist.broadcast广播给其他进程;或者所有进程各自加载同一文件(如果使用共享文件系统)。
state_dict = torch.load('model.pt', map_location=torch.device(f'cuda:{local_rank}'))
model.module.load_state_dict(state_dict)
特别注意:
-
必须使用
model.module访问原始模型:因为DDP包装了模型,model.state_dict()的键会多出module.前缀,而保存和加载通常期望原始键。 -
保存检查点需包含优化器状态:若要从断点恢复,必须同时保存
optimizer.state_dict()和epoch等信息,且通常每个进程都需要保存自己的优化器状态(因为各进程的优化器状态可能略有不同,尤其在混合精度训练中)。 -
多卡写入冲突:非rank 0进程不应保存模型权重,但优化器状态可以各存各的。
-
map_location:加载时需明确指定设备,以避免所有参数加载到rank 0的GPU上导致显存溢出。
完整保存恢复示例:
# 保存
if dist.get_rank() == 0:
torch.save({
'model': model.module.state_dict(),
'optimizer': optimizer.state_dict(),
'epoch': epoch,
}, 'checkpoint.pt')
# 恢复
checkpoint = torch.load('checkpoint.pt', map_location=torch.device(f'cuda:{local_rank}'))
model.module.load_state_dict(checkpoint['model'])
optimizer.load_state_dict(checkpoint['optimizer'])
start_epoch = checkpoint['epoch']
模型并行和张量并行有什么区别?¶
-
模型并行(Model Parallelism, MP) 是一个更宽泛的概念,指将整个神经网络的不同部分(如不同的层)划分到不同的设备上。流水线并行(Pipeline Parallelism, PP) 是其经典实现,按层切分,层与层之间传递激活值和梯度。
-
张量并行(Tensor Parallelism, TP) 是模型并行的一种特化,它不按层切分,而是将单个层内的权重矩阵(张量)沿着维度切分到多个设备上,每个设备计算该层的一部分,然后通过通信合并结果。它需要极高带宽的卡间互联(如NVLink),通信量极大,但能直接线性降低单卡权重和激活的显存。
核心区别:
| 维度 | 模型并行(PP) | 张量并行(TP) |
|---|---|---|
| 切分对象 | 模型的不同层(纵向) | 同一层内的权重矩阵(横向) |
| 通信模式 | 层间激活传递(数据量较小) | 层内高频All-Gather/All-Reduce(数据量极大) |
| 适用层级 | 跨节点、跨机器 | 同一节点内(依赖高速互联) |
| 显存节省 | 减少单卡层数,激活减小 | 减少单卡层内权重和激活占用 |
| 典型实现 | GPipe, 1F1B | Megatron-LM Tensor Parallelism |
在实际大模型训练中,两者常组合使用(3D并行),以最大化扩展效率。
如何使用PyTorch的弹性训练功能?torch.distributed.elastic。¶
弹性训练(Elastic Training)允许分布式训练任务在运行时动态适应节点数目的变化(例如因抢占式实例回收或节点故障导致节点减少,或集群扩容时节点增加),而无需从头开始训练。PyTorch提供了torch.distributed.elastic模块来支持这一功能。
核心组件:
-
torchrun(或torch.distributed.run):弹性启动器,它会自动设置环境变量,监控节点的变化,并管理进程的创建和销毁。 -
DynamicRendezvousHandler:负责动态节点发现和同步。它使用一个“会合点”(rendezvous)作为协调中心,允许节点加入和离开。 -
ElasticAgent:负责每个节点上的本地进程管理,包括启动、监控和重启工作进程。
使用步骤:
-
编写训练脚本:脚本中只需包含标准的分布式训练逻辑,无需处理节点变化。通过环境变量获取
WORLD_SIZE、RANK、LOCAL_RANK等信息。关键是要能从检查点恢复训练。 -
使用
torchrun启动:
torchrun --nnodes=1:4 --nproc_per_node=8 \
--max_restarts=3 --rdzv_id=job1 --rdzv_backend=c10d \
--rdzv_endpoint=localhost:12345 \
train.py
-
--nnodes:指定节点数的范围,如1:4表示最少1个节点,最多4个节点。 -
--rdzv_backend和--rdzv_endpoint:指定会合后端的类型和地址,用于节点发现和协调。 -
--max_restarts:进程组在失败后可重启的最大次数。 -
在训练脚本中处理弹性环境:
import torch.distributed.elastic.multiprocessing.errors as errors
@errors.record
def main():
# 标准分布式初始化
torch.distributed.init_process_group(backend='nccl')
# ... 训练循环 ...
if __name__ == "__main__":
main()
使用@errors.record装饰器可以捕获由于节点变化导致的进程错误,并打印出诊断信息。
- 支持动态成员变更:当节点数量发生变化时,
torchrun会自动杀死旧的工作进程,重新初始化分布式环境(新的world_size),并重新启动训练脚本。训练脚本需要能够从最新的检查点恢复,从而在新的节点规模下继续训练。
训练脚本内部的关键逻辑:
-
使用
DistributedSampler并设置正确的epoch和shuffle。 -
在每个epoch或固定步数后保存检查点,包括模型、优化器、学习率调度器和数据加载器的状态(如sampler的epoch)。
-
在启动时检查是否存在检查点,若存在则加载并恢复,同时根据当前的
world_size重新构建数据加载器和模型(若使用了模型并行需重新分区)。
优势:
-
容错性:节点失效时自动重启训练,减少人工干预。
-
成本优化:可利用抢占式实例,当实例被回收时自动迁移到新节点。
-
动态扩缩容:当集群有空闲资源时自动增加训练节点,提高训练速度。
注意事项:
-
弹性训练要求底层存储(如NFS、对象存储)对所有节点可达,以便共享检查点。
-
对于大规模模型,模型并行策略(如TP/PP)通常固定不变,弹性主要体现在数据并行的节点数变化。若需改变模型并行度,通常需要更复杂的重分片逻辑,PyTorch弹性目前主要支持数据并行的弹性。
在分布式训练中,如何处理变长数据的padding和对齐?¶
在分布式训练中,当使用DistributedSampler时,每个进程获得的是数据集的互不重叠子集。这些子集内部的样本长度仍然是变长的。处理变长序列的方法与单机训练类似,但需要注意分布式环境下的负载均衡和效率。
处理方法:
- 动态填充(Dynamic Padding):在每个batch内,将序列填充到该batch中最大长度,而不是整个数据集的最大长度。这通过自定义
collate_fn实现。
def collate_fn(batch):
inputs, labels = zip(*batch)
# 填充到batch内最大长度
padded_inputs = torch.nn.utils.rnn.pad_sequence(inputs, batch_first=True, padding_value=0)
labels = torch.tensor(labels)
return padded_inputs, labels
在DDP中,每个进程独立完成自己batch内的填充,无需跨进程通信。
-
序列打包(Packing):将多个短序列拼接成一个长序列,并生成对应的
attention_mask(分块对角掩码)。这可以在collate_fn中实现,也可以预先离线打包。打包后序列长度统一,提高了GPU利用率,且不需要padding。但需要正确处理attention_mask以隔离不同样本。 -
长度分桶(Length Bucketing):在构建
Sampler时,将长度相近的样本放入同一个batch,减少填充量。可以自定义BatchSampler,根据样本长度排序后分批。在分布式环境下,可以在每个进程内独立进行分桶,或者先在各进程内排序再采样。 -
负载均衡考虑:由于每个进程获取的数据子集是随机划分的,可能会出现某个进程的batch平均长度显著大于其他进程的情况,导致该进程计算变慢,成为“掉队者”(straggler)。为缓解此问题,可以在
DistributedSampler中使用shuffle=True,或者采用更精细的数据划分策略(如基于长度进行全局排序后再分配)。
对齐要求:对于某些模型架构(如需要pack_padded_sequence的RNN),需要按长度降序排列序列,并在collate_fn中同时返回排序索引和实际长度。现代Transformer通常使用attention_mask来忽略填充位置,对齐要求较低。
如何在DDP训练过程中动态改变batch size或序列长度?¶
在DDP训练过程中动态改变batch size或序列长度通常是出于显存优化或课程学习的目的,但需要谨慎处理,因为所有进程的数据采样必须保持同步。
方法一:梯度累积(等效动态batch size)
最安全且常用的方式。保持每个GPU上的micro-batch size固定,通过动态调整梯度累积步数来改变等效全局batch size。梯度累积在每步反向之间不进行All-Reduce,只在累积步数达到设定值时才同步梯度并更新参数。累积步数可以在训练过程中根据条件(如loss变化、epoch)动态改变。这种方式完全不需要修改数据加载器或分布式采样器,所有进程依然保持同步。
方法二:动态调整DataLoader的batch_size
如果直接修改DataLoader的batch_size,需要重新创建DataLoader。在DDP中,由于DistributedSampler已经将数据划分给各进程,新创建的DataLoader必须保证各进程划分的一致性。通常的做法是:
-
在每个epoch开始前,根据策略决定新的
batch_size。 -
重新创建
DistributedSampler(设置相同的epoch和seed)和DataLoader。由于种子一致,各进程仍然会获得不重叠的数据子集。 -
此方法较繁琐,且会导致每个epoch的数据加载逻辑重置。
方法三:动态改变序列长度
序列长度可以在数据预处理阶段通过截断(truncation)或padding来控制。在训练过程中动态改变最大序列长度,需要在DataLoader的collate_fn中实现动态截断逻辑。同样,由于每个batch的长度会变化,梯度累积步数和显存占用也会波动。一般不推荐在DDP中单独改变某个进程的序列长度策略,因为可能导致进程间负载严重不均。应采用全局统一的长度调整策略(如基于当前epoch逐渐增大长度)。
注意事项:
-
任何改变都必须在所有进程上同步执行,避免进程间逻辑不一致。
-
如果在训练过程中调整学习率调度器或其他超参,也应确保这些操作在所有进程上同时发生(通常只在rank 0上执行,然后broadcast或通过文件同步)。
7. 分布式训练报错“Address already in use”通常怎么解决?¶
该错误表示指定端口已被占用。
解决方法:
-
更换端口号:换一个未被占用的端口,例如将
MASTER_PORT从 29500 改为 12356。 -
杀掉占用进程:使用
lsof -i :端口号找出 PID,然后kill -9 PID。 -
设置随机或自动分配端口:在启动脚本中设置
MASTER_PORT为一个随机数。 -
检查并杀死残留进程:用
ps aux | grep python检查并杀死所有残留的训练进程。 -
确认
init_method配置正确:如果使用tcp://,确保 IP 正确且端口未被防火墙阻塞。使用torchrun的--rdzv_endpoint可自动管理端口。
18. 在多节点训练时,如何设置NCCL环境变量优化通信?¶
关键 NCCL 环境变量:
-
NCCL_DEBUG=INFO:开启调试信息,用于诊断性能瓶颈。 -
NCCL_SOCKET_IFNAME:指定网络接口名称(如eth0,bond0)。 -
NCCL_IB_DISABLE:无 InfiniBand 时设为1,强制使用 TCP/IP。 -
NCCL_NET_GDR_LEVEL=5:启用 GPU Direct RDMA(需硬件支持),大幅降低延迟。 -
NCCL_ALGO:AllReduce 算法,通常默认Ring。 -
NCCL_MIN_NCHANNELS:增加通道数可提高带宽利用率,但增加显存开销。 -
NCCL_CROSS_NIC=1:允许多网卡通信,提高带宽。
在作业脚本中导出这些变量即可,例如 export NCCL_SOCKET_IFNAME=eth0。
19. 解释“计算图不一致”错误在DDP中的典型原因及解决。¶
该错误表示不同进程上的模型执行了不同的前向计算路径,导致计算图结构不一致。
典型原因:
-
条件语句依赖输入数据:如
if x.size(1) > 100在不同进程进入不同分支。 -
某些参数只在部分进程中参与了前向计算,导致梯度集合不同。
-
数据加载逻辑导致 batch 大小不一致。
解决方法:
-
确保模型结构在所有进程上完全一致,避免依赖数据的控制流。若必须,使用张量操作(如
torch.where)代替if-else。 -
使用
DistributedSampler并设置drop_last=True保证每批大小一致。 -
临时方案:在 DDP 初始化时设置
find_unused_parameters=True,让 DDP 自动将未参与计算的参数梯度设为 0,从而保证集合一致(有额外开销)。
20. 如何在DDP中只同步部分参数(如不更新某些层)?¶
冻结某些层:将这些层的 requires_grad 设为 False。DDP 会自动检测哪些参数需要梯度,并仅同步这些参数的梯度。冻结的参数不会产生梯度,也不会参与 All-Reduce。
for param in model.backbone.parameters():
param.requires_grad = False
model = DDP(model, device_ids=[local_rank])
若希望参数需要梯度但独立更新(不同步),通常不推荐在 DDP 中这样做,因为这会导致模型不一致。可通过自定义通信钩子来跳过某些参数的同步,但实现复杂。
使用torch.distributed.algorithms.join处理进程掉队问题。¶
join 上下文管理器允许快速进程在等待掉队者时暂时“退出”并执行一些等待操作,防止死锁和超时。
from torch.distributed.algorithms.join import Join
with Join([model], enable=True):
for batch in dataloader:
# 训练循环
当某进程完成本地迭代而其他进程仍在工作时,Join 接管快速进程,使其进入等待循环,直到所有进程完成。它不能根除掉队问题,只是避免死锁,但总吞吐受限于最慢进程。更好的做法是优化数据加载和负载均衡。
FSDP与梯度检查点结合使用时有什么注意事项?¶
FSDP 和梯度检查点结合使用极为常见,两者都是节省显存的有效手段,配置时需注意:
-
分片粒度与检查点粒度协调:通常采用逐层(per-layer)FSDP 包裹,并对每个 Block 启用梯度检查点。这样前向完成后参数分片和激活几乎立即释放,显存节省最大化。
-
激活生命周期:两者在显存节省上是正交和累加的,检查点丢弃激活,FSDP 释放参数分片。
-
通信与计算重叠:反向时的重计算为 FSDP 的 All-Gather 和 Reduce-Scatter 提供更多计算时间,有助于隐藏通信延迟。
-
配置顺序:必须在 FSDP 包裹之前启用梯度检查点。FSDP 的
auto_wrap_policy应设为按层包裹。可调整forward_prefetch和backward_prefetch以平衡通信和计算。 -
与混合精度兼容,重计算的激活精度与 AMP 一致。
最佳实践:先启用 FSDP (ZeRO-3),再逐层检查点,最后叠加上 Flash Attention,按需调整预取。
分布式检查点如何保存和加载?与单卡检查点的区别。¶
分布式检查点:在大模型训练中,单卡无法容纳完整模型,需将模型参数和优化器状态分片保存在各进程的存储中。
-
保存:各进程调用
torch.save保存自己持有的参数分片、优化器状态分片及元数据。通常所有进程都保存,并依赖共享文件系统。FSDP 提供SHARDED_STATE_DICT高效保存分片,也可用FULL_STATE_DICT在 rank 0 聚合完整状态(需大量 CPU 内存)。 -
加载:训练恢复时,各进程从对应分片文件加载,要求当前分布式拓扑与保存时一致。FSDP 的
load_state_dict可自动处理分片映射。 -
与单卡检查点的区别:单卡检查点包含完整模型权重,可直接用于推理或迁移;分片检查点不是完整权重,依赖于分布式配置,无法直接用于单卡,但保存/加载极快且不消耗额外内存。可通过工具(如 DeepSpeed 的
ds_to_hf.sh)将分片转换为单卡检查点。
注意事项:保存时需所有进程同步(dist.barrier),确保文件写入完成。恢复时必须同时加载优化器状态和学习率调度器。
如果要训练一个百亿参数模型,你会如何设计分布式方案?¶
以训练 70B-175B 稠密大模型为例,设计一套典型的 3D 并行方案。
-
显存估算:70B 模型,FP16 权重 140 GB,梯度 140 GB,Adam 优化器状态 560 GB,总状态约 840 GB,单卡无法承受。
-
并行策略组合:
-
张量并行 (TP):在单机内使用,利用 NVLink 高带宽。将单层内矩阵沿隐藏维切分,TP 度通常设为单机 GPU 数(如 8),单卡权重/激活显存降为 1/8。
-
流水线并行 (PP):跨节点使用,将模型按层切分到不同节点,层间传递激活,通信量小。PP 深度可根据模型层数和节点数设定,例如 PP=4 或 8。配合 micro-batching 和 1F1B 调度减少气泡。
-
数据并行 (DP) / ZeRO-3:在更大范围复制 TP×PP 模型副本,处理不同数据。ZeRO-3 进一步分片优化器和参数到所有 DP 进程,最大化显存节省。
配置实例 (70B 模型,256 卡 A100,32 节点):
-
TP=8 (每节点 8 卡,节点内)
-
PP=4 (跨 4 节点)
-
DP = 总卡数 / (TP × PP) = 256 / (8×4) = 8,即 8 个数据并行副本。
-
结合 ZeRO-3 在 DP 组内分片参数和优化器。
-
显存优化:开启 BF16 混合精度,逐层梯度检查点,Flash Attention-2。若仍紧张,将优化器状态 Offload 到 CPU (ZeRO-Offload)。
-
数据加载:使用
DistributedSampler,离线预处理为 tokenized 格式,支持流式加载。 -
容错与监控:
torchrun弹性训练,定时保存分片检查点,部署 NCCL 和 GPU 监控告警。 -
训练流程:采用余弦退火学习率 + warmup,利用 DeepSpeed 或 Megatron-LM 等成熟框架实现。该方案充分利用硬件拓扑,在显存、通信和计算之间取得平衡。