跳转至

分布式训练

数据并行(DP 与 DDP)的原理是什么?PyTorch DDP 是如何实现梯度同步的?

image.png

image.png

数据并行(Data Parallelism) 是最常用的分布式训练策略,其核心思想是将一个 batch 的数据分割成多个 micro-batch,分别在多个设备(GPU)上各自持有一份完整的模型副本,独立完成前向传播和反向传播得到梯度,然后将所有设备上的梯度进行汇总平均,用平均梯度更新每个副本的模型参数,从而保持所有副本的一致性。

DP(DataParallel,单机多卡)

  • PyTorch 的 nn.DataParallel 是一种单进程多线程的实现。

  • 工作流程:主 GPU(通常为 GPU0)将 batch 切分并广播给其他 GPU;所有 GPU 完成前向计算;各 GPU 的输出 gather 到主 GPU 计算损失;损失广播到各 GPU 进行反向传播;各 GPU 的梯度在主 GPU 上进行求和(reduce),然后更新主 GPU 的参数,再将参数广播给其他 GPU。

  • 缺点:主 GPU 负载不均衡(需聚合输出、损失、梯度,再广播参数),通信效率低;Python 全局解释器锁(GIL)限制多线程性能;仅支持单机。

DDP(DistributedDataParallel,多机多卡)

  • 采用多进程机制,每个 GPU 由一个独立进程控制,避免了 GIL。

  • 模型副本:每个进程初始化相同的模型,加载相同的初始参数。

  • 前向与反向:每个进程独立执行前向和反向,获得各自 local batch 上的梯度。

  • 梯度同步:通过AllReduce通信原语,将各进程的梯度进行求和(或平均),使所有进程获得完全相同的平均梯度。同步发生在每次反向传播计算出所有梯度之后、优化器 step 之前。

  • 参数更新:每个进程用相同的平均梯度更新自己的模型,保证所有副本在每次迭代后参数严格一致。

  • DDP 还提供了 gradient_as_bucket_view 和通信与计算 overlap 的优化:将梯度分桶,一旦某个桶内的梯度全部计算完毕,立即启动该桶的 AllReduce,同时继续计算其他层的梯度,从而隐藏通信时间。

梯度同步实现(AllReduce) DDP 使用 NCCL 后端的 AllReduce 操作,支持多种算法(如 Ring AllReduce、Tree AllReduce)。它对所有参数的梯度张量进行规约(默认求和),然后除以进程数实现平均。DDP 通过 Reducer 机制在 autograd 梯度产生时就异步触发通信,实现高效同步。


DDP 中的 AllReduce 梯度平均是如何进行的?解释 Ring AllReduce 的通信过程。

DDP 中的梯度平均是通过 AllReduce 操作完成的。在训练中,每个 GPU 计算出的梯度是局部小批量的梯度,需要将所有 GPU 的梯度求和并取平均,以获得全局 batch 的梯度。AllReduce 就是所有进程各自提供数据,经过规约(求和)后,所有进程都得到相同的结果。DDP 随后将结果除以 world_size 得到平均梯度。

Ring AllReduce 是实现 AllReduce 的一种高效算法,尤其适用于大规模分布式系统。其通信过程分为两个阶段:Scatter-Reduce 和 AllGather。假设有 N 个 GPU 排成一个逻辑环,每个 GPU 的待规约数据(如梯度)被切分为 N 个块。

  • 第一阶段:Scatter-Reduce 进行 N−1 次迭代。每次迭代中,每个 GPU 将本地的一个数据块发送给环上的下一个 GPU,同时从上一个 GPU 接收一个数据块,并与自己的对应块累加。经过 N−1 次后,每个 GPU 上都持有一块完整的规约结果(即该块在所有 GPU 上的和)。此时,梯度块分散在不同的 GPU 上。

  • 第二阶段:AllGather 同样进行 N−1 次迭代。每个 GPU 将自己在第一阶段结束时持有的完整规约块发送给下一个 GPU,并接收上一个 GPU 传来的另一个已规约块。最终,每个 GPU 都收集到所有块的规约结果,即完整的全局梯度。

特点:Ring AllReduce 的通信量每个 GPU 为 2(N−1)/N×总数据量2(N−1)/N×总数据量,带宽利用率高,且每个 GPU 的通信负载均衡,是 DDP 默认通信算法之一。

DDP 通过 torch.distributed.all_reduce 调用 NCCL 实现这一操作,并可以对梯度张量进行分桶,实现通信与反向传播的并行。


模型并行是在什么情况下必须使用的?它与数据并行如何区分?

模型并行(Model Parallelism) 是指将模型的参数和计算图切分到不同的设备上,当单个设备无法容纳整个模型时,必须使用模型并行。例如,训练巨大的 Transformer 模型(数百亿甚至万亿参数),单卡显存远不够存放模型权重和中间激活,此时必须将模型的不同部分放置在不同的 GPU 上。

与数据并行的区分

与数据并行的区分

  • 数据并行:每个设备拥有完整模型副本,数据切分到各设备独立计算,同步梯度。要求每个设备能装下完整模型。

  • 模型并行:模型被切分,每个设备只负责一部分参数和计算,设备间需要传递中间激活和梯度。不同设备合作完成一个 batch 的计算。数据可以是一个 batch 的不同 micro-batch,也可以是对同一数据的不同层。当单卡无法放下模型时,模型并行成为必需。

进一步细分

  • 流水线并行(按层切分)

  • 张量并行(按层内矩阵乘法切分)

  • 序列并行(按序列长度切分)

实践中常将数据并行与模型并行混合使用,例如 Megatron-LM 使用张量并行 + 流水线并行 + 数据并行的 3D 并行策略来训练超大规模模型。


流水线并行(Pipeline Parallelism)如何划分模型?GPipe 和 PipeDream 的调度策略有何不同?

流水线并行 将模型按层(或层组)垂直切分,每个设备(GPU)负责若干连续的层,形成流水线。一个 batch 被进一步划分为多个 micro-batch,以流水线方式送入第一个设备,后续设备依次处理,使得设备利用率提高。

GPipe 调度策略

  • 将 micro-batch 顺序注入流水线,所有 micro-batch 的前向传播完成后,再进行反向传播。

  • 由于前向和反向分离,需要暂存所有 micro-batch 的激活,直到反向时使用,因此显存占用较大。

  • 优点:实现简单,流水线气泡(idle time)存在但可通过增大 micro-batch 数量来降低气泡比例。

  • 周期性注入 micro-batch,形成“前向-等待-反向”的节奏,但等待时间(bubble)随 micro-batch 数量增加而减少。

PipeDream 调度策略

  • 使用1F1B(one forward one backward) 的交替调度:一旦某个设备完成一个 micro-batch 的前向,立即开始另一个 micro-batch 的反向(而不是等所有前向完成)。

  • 这种交错可以更早地释放激活占用的显存,降低显存峰值。

  • 需要更复杂的调度逻辑和版本管理,但整体气泡更小,效率更高。

  • PipeDream 还需要考虑参数版本一致性(因为前向和反向可能使用不同版本的参数),通常配合权重更新点的控制。

区别总结:

GPipe 简单,气泡可控(通过微批数量),但显存高;PipeDream 更高效,显存占用更低,但实现复杂度高。


张量并行(Tensor Parallelism)是如何拆分 Transformer 层的矩阵乘法的?

张量并行将单个层内的矩阵乘法(如全连接层、多头注意力中的投影)沿着参数维度切分到多个 GPU,实现单层内的并行。以 Transformer 的 FFN 和 Attention 为例,Megatron-LM 提出的方案:

MLP 的拆分

image.png

自注意力的拆分

注意力头的数量天然适合并行:每个 GPU 负责一部分头(head)的 Q、K、V 投影和注意力计算,然后各自计算得到部分输出,再通过 AllReduce 汇总。对于 Q、K、V 的权重矩阵,可以按列切分,每个 GPU 计算自己的头。

通信:在 MLP 的切分中,前向传播中间需要一次 AllReduce 求和(或后向时相应通信)。这种张量并行在单机多卡内通信开销小(NVLink 高速互联),常与流水线并行结合使用。


什么是序列并行(Sequence Parallelism)?它通常与哪种并行方式组合?

序列并行 是将输入序列的长度维度切分到不同设备上,每个设备负责序列的一部分 token。这对于极长序列(如训练超长上下文模型)可以减少单个设备上的激活显存。通常与张量并行或流水线并行结合使用。

典型组合:

  • 在 Megatron-LM 中,序列并行与张量并行结合:在张量并行中,LayerNorm 和 Dropout 等操作沿序列维度独立,为了减少冗余显存,将这些操作也进行序列切分。在 Transformer 层中,LayerNorm 和残差连接在序列维度上是并行的,可以将激活在序列维度拆分,然后仅在需要张量并行的矩阵乘法(如 attention 或 FFN)之前进行集合通信转换。

  • 具体实现:输入序列的 token 被划分到多个设备,每个设备上执行部分的 LayerNorm、Dropout,而在执行张量并行的矩阵乘法时,通过 all-gather 收集完整序列,或者保持拆分并使用 reduce-scatter 来更新。

  • 序列并行能够大幅度减少激活的内存占用(因为这些操作本身没有跨序列依赖),与张量并行结合能训练更长的序列。

优势:扩展模型支持的序列长度,降低显存峰值。


ZeRO(DeepSpeed)的三个阶段分别对哪些状态(参数、梯度、优化器状态)进行分片?各有什么通信开销?

ZeRO(Zero Redundancy Optimizer)将优化器状态、梯度和参数在数据并行的设备间进行分片,消除冗余存储,从而支持更大模型。

ZeRO-1 (Optimizer State Sharding)

  • 分片内容:优化器状态(如 Adam 中的动量、方差)。

  • 参数和梯度:仍完整存储在每个 GPU 上。

  • 通信开销:在优化器更新步骤中,每个 GPU 需要收集所有其他 GPU 上的分片优化器状态(或收集更新后的参数),涉及 AllGather 或 ReduceScatter,通信量与数据并行相当(与模型参数量成比例)。

ZeRO-2 (Optimizer State + Gradient Sharding)

  • 分片内容:优化器状态 + 梯度分片。

  • 参数:仍完整存储。

  • 通信开销:每个 GPU 只存储与其优化器状态分区对应的梯度。在反向传播后,需要将梯度通过 ReduceScatter 规约并分片,各 GPU 只保留自己负责的部分梯度。通信量仍与模型大小同阶,但避免了冗余存储梯度。

ZeRO-3 (Optimizer State + Gradient + Parameter Sharding)

  • 分片内容:优化器状态、梯度、参数全部分片。每个 GPU 只存储模型参数的 1/N。

  • 前向/反向传播时:需要动态收集所需的参数。通信开销更高:前向时需要 AllGather 收集参数,反向时需要再次 AllGather 收集参数,然后 ReduceScatter 梯度。总通信量为标准数据并行的 1.5 倍左右,但将显存占用降到了 1/N,能够训练巨型模型。

总结:ZeRO 的阶段越高,显存节省越多,但通信量从不变到略微增加(ZeRO-3 通信量约为 DP 的 1.5 倍),需要在模型规模和通信开销间权衡。


ZeRO-Offload 是如何将优化器状态和计算卸载到 CPU 的?会带来多大速度损失?

ZeRO-Offload 将优化器状态(以及梯度、参数,取决于 ZeRO 阶段)从 GPU 显存卸载到 CPU 内存,同时将一部分计算(如优化器的更新步骤)也迁移到 CPU 执行,从而极大地降低 GPU 显存占用,使得在有限 GPU 资源下训练巨大模型成为可能。

工作流程:

  • 前向和反向传播仍在 GPU 上进行(计算密集部分)。

  • 反向传播得到的梯度被传输到 CPU。

  • 在 CPU 上执行优化器更新(如 Adam 的动量更新、权重衰减),更新参数。

  • 更新后的参数传输回 GPU 用于下一轮迭代。

  • 可以采用 CPU 与 GPU 计算重叠的策略(如 DeepSpeed 的 CPU-Adam),在 GPU 执行下一 micro-batch 前向的同时,CPU 进行优化器更新和参数传输,以隐藏部分延迟。

速度损失:

  • 受限于 PCIe 带宽和 CPU 计算速度。如果模型极大,GPU 计算成为瓶颈,卸载带来的通信延迟可以被计算掩盖,整体吞吐量损失较小(可能 10-20% 或更低)。但如果模型相对较小,GPU 计算很快,则 CPU 卸载和传输可能成为瓶颈,吞吐量可能下降数倍。

  • 实际效果依赖于硬件配置、模型大小、batch size 等。DeepSpeed 通过高效的流和缓存优化,尽力将开销降到最低。

ZeRO-Offload 结合 ZeRO-3 可以实现用单张或少数几张 GPU 训练超大模型。


混合精度训练的基本流程是什么?为什么需要 FP32 的主权重副本?

混合精度训练(Mixed Precision Training)同时使用 FP16(半精度浮点)和 FP32(单精度浮点)来加速训练并降低显存,同时保持模型精度。

基本流程:

  1. 维护一份 FP32 的主权重副本(master weights)。

  2. 每次迭代,将 FP32 权重转换为 FP16,得到 FP16 模型副本。

  3. 前向传播:使用 FP16 权重和 FP16 输入进行计算(矩阵乘法等利用 Tensor Core 加速)。

  4. 计算损失:通常仍为 FP32,以便更精确。

  5. 反向传播:在 FP16 下计算梯度(部分操作可转 FP32 提升稳定性)。

  6. 梯度转为 FP32,并乘以损失缩放因子(见下一问)以恢复被压到零的梯度值。

  7. 使用 FP32 梯度更新 FP32 主权重副本。

  8. 重复。

为什么需要 FP32 主权重:

image.png

  • 在 FP32 中累积权重更新可以保持足够的精度,避免舍入误差累积。

  • 此外,某些操作(如 BN、Softmax、大规约)在 FP32 下更稳定,防止数值问题。

因此,FP16 用于计算加速,FP32 用于准确存储和更新权重。


FP16 训练中,损失缩放(Loss Scaling)的原理是什么?动态缩放如何工作?

原理

image.png

静态缩放:选择一个固定的缩放因子,通常根据经验或尝试,但可能不适配所有训练阶段。

动态缩放

在训练过程中自动调整缩放因子。基本算法:

image.png

  1. 每个训练迭代,进行混合精度的前向与反向,梯度放大后更新。

  2. 检查梯度是否出现溢出(inf 或 NaN)。如果没有溢出,则可能逐渐增大缩放因子(如乘以 1.05),以提高对微小梯度的捕获。

  3. 如果检测到溢出,则该次更新跳过,缩放因子减小(如除以 2),避免后续梯度继续溢出。

  4. 通过这种自动调整,在训练大部分时间保持较大缩放因子,同时防止崩溃。

PyTorch 中 torch.cuda.amp.GradScaler 实现了动态缩放,广泛使用。


BF16 与 FP16 相比,为什么不需要损失缩放?在什么场景下更适用?

BF16(Brain Floating Point 16) 具有与 FP32 相同的 8 位指数位,但只有 7 位尾数;而 FP16 有 5 位指数、10 位尾数。

为什么 BF16 不需要损失缩放

image.png

  • BF16 的缺点是尾数精度更低(7 位 vs 10 位),表示数值的粒度更粗糙。但实践中神经网络对权重和梯度的精细度不敏感,BF16 通常能保持与 FP32 相当的模型精度。

适用场景

  • BF16 需要硬件支持(如 NVIDIA A100、Google TPU、Intel Sapphire Rapids)。在这些设备上,BF16 可以直接替代 FP32 进行训练而无需损失缩放,简化训练流程,且因为尾数更短计算效率更高。

  • 适用于大规模训练,如大模型、推荐系统等,显存带宽节省和计算加速效果显著。

  • 对于精度极度敏感的任务(如某些科学计算),BF16 可能精度不足,需混合使用。

对比:FP16 需损失缩放但尾数精度高;BF16 免缩放,动态范围大,更易用。


在分布式训练中,如何保证不同 GPU 上的 Batch Normalization 同步?(SyncBN)

标准的 BN 层在每个设备上独立计算当前 micro-batch 的均值和方差。当数据并行时,每个 GPU 上的 mini-batch 较小,统计量估计不准确,导致模型精度下降。SyncBN(Synchronized Batch Normalization) 跨 GPU 同步 BN 的统计量,使得每层的均值和方差是基于整个 global batch 计算的。

实现方式:

  • 在各设备完成前向传播经过 BN 层时,收集各设备上该层的局部均值和局部方差(以及样本数)。

  • 使用 AllReduce 或其他通信操作,计算全局均值和方差:

  • 全局均值 = 所有设备均值的加权平均(按样本数)。
  • 全局方差 = 根据全局均值和各设备局部方差、局部均值计算得到。

  • 然后每个设备用全局统计量对该设备上的特征进行归一化,并继续前向传播。

  • 反向传播时,梯度同样需要同步,确保每台设备上的 BN 参数(γ,β)梯度一致。

PyTorch 支持:torch.nn.SyncBatchNorm 封装了上述逻辑,可无缝替换普通 BN。在使用 DDP 时,只需将模型中的 BN 层转换为 SyncBN,然后 DDP 会自动处理所需的 AllReduce 通信,而不需额外代码。注意 SyncBN 会增加通信开销,但能有效提升小 batch size 分布式训练的精度。

注意事项:SyncBN 在所有设备上必须统一,确保每个进程输入相同的统计量计算;对于小 batch size 场景尤其重要;在混合精度训练中,SyncBN 一般在 FP32 下执行以保持精度。


大规模训练时,梯度累积(Gradient Accumulation)的作用是什么?它会影响 BN 吗?

作用

梯度累积是指在多个连续的 micro-batch 上计算梯度,但不立即更新参数,而是累加这些梯度,直到累积到等效于目标大 batch size 的数量后,再执行一次参数更新。这样可以在显存有限的条件下模拟大 batch 训练。例如,目标 global batch size 为 256,但单卡每次只能处理 32,则累积 8 步后更新,相当于 batch size 256。它几乎不增加显存,却可以增大有效 batch size,稳定训练,并充分利用 GPU 的算力。

对 Batch Normalization 的影响

标准的 BN 在每个 micro-batch 内部独立计算均值和方差,不跨累积步。因此,梯度累积会导致 BN 统计量不准确:BN 看到的只是小 batch(如32),而非全局 batch(256)的统计,这会降低模型的性能。解决方案:

  • 使用 SyncBN(跨卡同步 BN 统计量),但仍然不能解决累积步之间统计独立的问题。SyncBN 可以保证每个卡上的 batch 统计被同步,但累积步之间 BN 的均值和方差仍然来自于各个 micro-batch 的局部统计。

  • 更准确的做法是计算跨累积步的全局统计量:在训练时,可以利用所有累积步的均值和方差重新计算公式,或者使用 Ghost BN、Virtual BN 等技巧,但往往实现复杂。实践中,如果累积步数不多,BN 影响可以接受,也可以使用 Group Normalization (GN) 或 Layer Normalization (LN) 等不依赖 batch 统计的归一化方法彻底避免该问题。


如何计算多机多卡训练的理论最大吞吐量?通信带宽如何影响?

理论最大吞吐量通常用每秒处理样本数(或 token 数)衡量,受限于单步计算时间与通信时间。

计算模型

image.png

通信带宽的影响

通信带宽直接决定了梯度同步耗时。当模型梯度很大、GPU 数量多时,通信可能成为瓶颈。此时吞吐量受限于带宽,即使计算再快,也无法超越通信墙。为了提升吞吐,需要高带宽互联(如 NVLink、InfiniBand),并使用通信计算重叠、梯度压缩等技术。

估算示例:若模型梯度总大小 10 GB,使用 100 Gbps 带宽(12.5 GB/s),则 AllReduce 每卡通信时间 ≈ 2*(10/12.5) ≈ 1.6 秒(粗略)。若计算时间 0.5 秒,则整体 iteration 受限于通信,吞吐量由带宽决定。


解释“梯度通信与计算重叠”如何利用 CUDA Stream 或异步通信。

在现代分布式训练框架(如 PyTorch DDP, Horovod, Megatron-LM)中,可以利用反向传播的顺序性,将梯度的计算与已经计算出的梯度通信进行重叠,隐藏通信延迟。

实现原理

  • DDP 内部会将模型参数的梯度分为多个“桶”(buckets)。在反向传播过程中,一旦某个桶内的所有参数的梯度都计算完毕,DDP 就立即启动该桶的 AllReduce 异步通信(使用 NCCL 的异步操作),同时反向传播继续计算后续参数的梯度,不需要等待该通信完成。

  • 这一切依赖于 CUDA Stream:DDP 创建一个独立的 CUDA Stream 用于通信,使得通信与计算 Stream 可以并行执行。在同一个 Stream 上,通信操作被正确地插入,不会与计算发生依赖冲突。

  • 最后,在优化器步骤开始前,需要确保所有桶的通信都已完成(通过同步事件)。

异步通信 NCCL 提供了非阻塞的集合通信 API(如 ncclAllReduce 配合 cudaStream_t),PyTorch 的 all_reduce 也可以设置 async_op=True,返回一个 work handle,随后可等待完成。DDP 通过管理这些异步 handle,实现最大限度地重叠。

image.png


当集群中有节点故障时,如何实现容错训练和自动恢复?

在大规模长时间训练中,节点故障不可避免。容错机制包括:

检查点(Checkpointing)

定期保存模型状态、优化器状态、随机数种子、学习率调度器状态等到持久存储(如共享文件系统或对象存储)。发生故障后,可以回滚到最近的检查点继续训练。这需要检查点保存和加载足够快速,且不会过度干扰训练吞吐。

心跳检测与自动重启

  • 使用分布式训练框架(如 Kubernetes + PyTorch Operator, Ray Train, DeepSpeed)监控每个 worker 的健康状态(心跳)。

  • 当某个 worker 崩溃或无响应时,框架可以自动终止当前训练,并利用之前保存的检查点在新的健康节点集合上重新启动作业。

  • 需要作业管理层面支持动态资源分配和作业重新提交。

弹性训练

更高级的是弹性训练(见下问),允许在训练过程中动态增减 worker,当节点故障时,可以从剩余健康节点继续训练,无需完全停止重启。Horovod 的弹性训练功能、PyTorch Elastic 等提供了这种能力。

TorchElastic PyTorch 提供了 torch.distributed.elastic 模块,可以监控 worker 状态,并在 worker 失败时重新初始化进程组,加载最近检查点继续训练。它能处理少数节点故障,保留大多数节点继续工作。

总结:通常采用定期检查点 + 弹性/自动重启,实现分钟级的恢复。


什么是弹性训练(Elastic Training)?它如何动态调整 worker 数量?

弹性训练允许分布式训练作业在运行时动态地增减计算资源(worker 数量),而无需从零重新启动训练。当新增节点时,能充分利用新增算力;当节点失效或需要释放资源时,可以在剩余节点上继续训练。

动态调整机制

  1. 资源管理:通常基于 Kubernetes 或 YARN,作业可以使用弹性调度器,根据可用资源自动调整 worker 数量。

  2. 重新初始化分布式进程组:当 worker 数量变化时,需要重新建立通信组。PyTorch Elastic 通过 rendezvous 机制,worker 们通过 etcd 或文件系统同步,等待所有参与节点都达到同一个状态,然后重新初始化 ProcessGroup

  3. 数据划分调整:数据加载器需要根据新的 worker 数量重新划分数据分片,以避免数据重复或遗漏。通常每个 worker 基于自己新的 rank 和 world size 重新计算分片。

  4. 模型与优化器状态:所有 worker 需要从最近的检查点恢复模型和优化器状态。因为检查点包含 world size 信息,弹性训练需要能够将检查点中的状态适配到新的 worker 数量(例如,调整批次归一化统计、数据加载器的 epoch 等)。

  5. 学习率与 batch size 调整:worker 数量变化会影响有效 global batch size,因此学习率也需要相应调整(如线性缩放)以保持训练动态。

弹性训练使得大规模训练能容忍机器池的波动,极大提升资源利用率和稳定性。


在 PyTorch 中,DistributedDataParallel 与 DataParallel 相比有何优势?

查看内嵌表格

总结:DDP 在性能和扩展性上全面优于 DP,已成为 PyTorch 分布式训练的标准。


初始化分布式进程组时,NCCL 和 GLOO 后端分别适用于什么?

  • NCCL (NVIDIA Collective Communications Library):专为 NVIDIA GPU 设计,支持极高性能的集合通信(AllReduce, AllGather, Broadcast 等)。适用于GPU 间通信,包括单机多卡和多机多卡。强烈推荐用于 GPU 训练,因为它优化了 GPU 拓扑和 NVLink/IB 网络。

  • GLOO:一个通用的集合通信库,支持 CPU 和 GPU。相对于 NCCL,GPU 上的性能较低,但可用于 CPU 分布式训练,或在不支持 NCCL 的环境中(如非 NVIDIA GPU 或某些特定网络)使用。也可以作为 GPU 训练的回退选择。

选择原则:

  • GPU 训练:优先 NCCL。

  • CPU 训练或异构环境:使用 GLOO。

  • 有时某些操作(如 send/recvbarrier)在 NCCL 上可能受限,GLOO 作为补充。


如何调试分布式训练中的死锁或挂起问题?

分布式训练中的死锁通常是由于集合通信操作不匹配或某个进程因错误退出,导致其他进程一直等待。调试方法包括:

  • 使用 NCCL_DEBUG 环境变量:设置 NCCL_DEBUG=INFOWARN,输出 NCCL 的详细日志,查看哪个 rank 通信卡住。

  • 检查进程存活:使用 torch.distributed.barrier() 前后添加日志,判断卡在哪个 barrier。也可以使用 gdb 附着到挂起的进程,查看调用栈。

  • 添加超时:在初始化进程组时设置 timeout 参数(如 init_process_group(backend='nccl', timeout=datetime.timedelta(minutes=10))),超时会抛出异常,避免无限挂起。

  • 确保所有 rank 执行相同的集合操作序列:死锁的常见原因是一个 rank 调用 all_reduce 而另一个没有,导致不匹配。检查条件分支,确保所有 rank 在相同位置调用相同的通信函数。

  • 数据加载器不一致:如果某些 rank 的数据加载器抛出异常并退出,而其他 rank 还在等待通信,会导致挂起。可以捕获异常并在所有 rank 上同步。

  • 使用 torch.distributed.monitored_barrier(PyTorch 1.10+)检测哪个 rank 未到达 barrier。

  • 日志记录:在每个关键步骤打印 rank 和时间戳,定位到哪一步停止。

通过这些手段,一般可以定位到具体的死锁位置。


使用梯度检查点(Gradient Checkpointing)时,分布式训练的显存如何进一步节约?

梯度检查点(Activation Checkpointing)是一种以计算换内存的技术。在反向传播期间,通常需要保存前向传播中的中间激活,以计算梯度。检查点技术只保存部分激活,其他在前向传播时丢弃;在反向传播需要时,通过重新计算一段小图来恢复这些激活。

在分布式训练中,它与张量并行、流水线并行结合使用:

  • 张量并行:每个设备只保存自己那部分激活,显存随设备数线性减小。结合检查点,可以进一步显著减少每卡的激活内存。

  • 流水线并行:1F1B 调度本身降低了激活存储,使用检查点后,微批的激活存储更少,可以支持更大的 micro-batch 或更深的模型。

  • 选择性检查点:只对计算量较大的层(如注意力、FFN)进行重计算,对便宜的操作保留激活。

通过这种组合,即使不增加设备,也能训练更大模型。梯度检查点在分布式训练中并不直接减少通信,而是解放显存用于更大的 batch 或模型,间接提升训练吞吐。


在 3D 并行中,如何确定最优的并行度配置(dp, tp, pp)?

3D 并行结合数据并行(DP)、张量并行(TP)和流水线并行(PP)。最优配置需根据模型结构、硬件拓扑、通信带宽、显存限制来权衡。

步骤与原则:

  1. 张量并行(TP):通常限制在单机内部,利用高带宽 NVLink(如 8 卡机)。TP 度一般等于单机 GPU 数量(如 8),超过单机会有极大通信开销,不划算。对于 Transformer,Megatron 设置 TP=8 常见。

  2. 流水线并行(PP):将模型层切分到多个设备。PP 的通信量相对较小(仅层边界激活和梯度),跨节点也可以接受。PP 度越大,流水线气泡越多,需要大量 micro-batch 来摊销气泡,因此总 batch size 要足够大。同时 PP 的延迟增加。

  3. 数据并行(DP):对剩余的维度使用数据并行,DP 通信(梯度 AllReduce)量与模型总参数成正比,通常需要高带宽网络。在 3D 中,由于 TP 和 PP 已经减少了每 GPU 的模型部分,DP 部分通信量相应降低。

确定过程:

  • 根据模型总参数和单卡显存,首先确定必须的模型并行度(TP×PP),确保每卡可容纳模型。

  • 在满足显存前提下,优先增大 DP 以提高吞吐,因为 DP 扩展性最好。

  • TP 保持在一个节点内(如 8),PP 可以根据模型层数灵活设置,通常 2~8。

  • 尽量让 TP×PP×DP = 总 GPU 数,并使得每个并行维度的通信开销与计算平衡。

  • 可以用分析工具(如 Megatron-LM 的计算器)模拟不同配置的吞吐量、显存和 bubble 率,选择最佳。

总结:TP 在内,PP 切深度,DP 扩规模,寻找全局最优。


分布式训练中,学习率如何根据 batch size 缩放?线性缩放规则和平方根缩放。

当 global batch size 增大时,需要相应调整学习率以保持训练稳定和收敛。

  • 线性缩放规则(Linear Scaling Rule):当 batch size 增大 k 倍时,学习率也线性增大 k 倍。理论基础:在 SGD 中,每一步梯度是基于 k 倍多样本的均值,方差减小为 1/k,因此增大学习率 k 倍可以保持参数更新量的方差大致不变。经验上,对于大批量训练,线性缩放常配合 warmup 使用,在 batch size 不是极端大时效果良好。

image.png

当前最佳实践(如大批量训练 BERT, GPT 等):使用线性缩放规则,但结合学习率 warmup(前几千步从零线性增加到目标学习率),然后再衰减。有时也会根据损失形状调整,或者使用 LARS 等自适应优化器以适应大 batch。最终缩放策略可以通过在小批量上调优后,在 scale 时调整得到。

image.png


如何在一个大型集群上高效地执行超参数搜索(Hyperparameter Sweep)?

在大规模集群上执行超参数搜索,需要高效的编排和资源利用策略。

常见方法:

  1. 独立作业并行:每个超参数配置作为一个独立的训练作业提交。适合作业调度系统(如 Slurm, Kubernetes)。缺点是每个作业都需要从头加载数据、初始化模型,且资源利用率可能不均衡。

  2. 分布式搜索框架:如 Ray Tune、Optuna、Hyperopt 结合分布式后端。这些框架可以启动多个训练 trial,每个 trial 占用一部分资源(如一个多节点组),动态管理 trial 的生命周期。支持早期停止(如 ASHA 算法)以终止表现不佳的 trial,节省资源。

  3. 多 worker 共享模型初始化:在某些场景中,可以一次预训练模型,然后只调整少量参数微调,减少重复计算。但在完全的搜索中不易做到。

高效策略:

  • 群体训练(Population Based Training, PBT):同时训练一组模型,定期利用表现好的模型的权重和超参数替换差的,实现自适应搜索。适合在线进化。

  • 贝叶斯优化:利用代理模型选择下一个尝试的超参数,减少搜索次数。

  • 使用 spot 实例/节点,配合容错检查点,降低计算成本。

  • 数据与模型的缓存:在多个 trial 间共享公共数据,或者使用内存缓存加速 I/O。

  • 资源粒度:不是所有 trial 都需要全数 GPU,可根据模型大小动态分配资源。

通过使用 Ray Tune + PyTorch Distributed,可以方便地在集群上启动分布式 trial,并与集群调度器集成,实现高效超参搜索。