跳转至

分布式训练并行策略:从数据并行到3D并行及通信优化

数据并行(DP vs DDP)的原理是什么?DDP 的 Ring AllReduce 是如何工作的?

数据并行(Data Parallelism) 是最基础的分布式训练策略。其核心思想是:每张 GPU 上保存一份完整的模型副本,训练数据被均匀划分到各 GPU。每个 GPU 使用自己的数据子集独立进行前向和反向传播计算,得到各自的梯度。然后所有 GPU 需要将这些梯度进行全局汇总(通常求平均),以保证模型副本之间保持一致。根据梯度同步方式的不同,分为 DP(Data Parallel)和 DDP(Distributed Data Parallel)。

DP(通常基于参数服务器):一台 GPU 作为参数服务器,其他为工作节点。工作节点各自计算梯度,发送给参数服务器;服务器聚合梯度、更新参数,再将新参数广播给所有工作节点。这种方案存在通信瓶颈和单点故障,且扩展性差。

DDP(分布式数据并行,Ring AllReduce):目前主流方案。它摒弃了中心服务器,所有节点形成一个逻辑环。梯度同步采用 Ring AllReduce 算法,整个过程不需要中心节点,通信负载均衡。

Ring AllReduce 工作流程:

  • 假设共有 NN 个 GPU,每个 GPU 上的梯度张量被划分为 NN 个连续的数据块。

  • Scatter-Reduce 阶段:执行 N−1N−1 步。在第 ii 步,每个 GPU 将本地一个数据块发送给下一个 GPU,同时从前一个 GPU 接收一个数据块。接收后,将该块与自己对应的块进行累加(reduce)。这样经过 N−1N−1 步,每个 GPU 都拥有一个完整归约好的数据块(即该块在所有 GPU 上的和)。

  • AllGather 阶段:同样 N−1N−1 步。每个 GPU 将已归约好的数据块依次传递给下一个 GPU,其他 GPU 存储该块但不累加。经过 N−1N−1 步,所有 GPU 都收集到了所有归约好的数据块,即获得了全局梯度求和结果。

最后,各 GPU 使用相同的全局梯度来更新自己的模型参数。DDP 的通信量是 2(N−1)2(N−1) 次数据块传输,总数据量为 2×梯度总大小×(N−1)/N,接近最优。Ring AllReduce 充分利用了每个节点的双向带宽,避免了中心服务器的瓶颈。

张量并行(Tensor Parallelism)是如何拆分 Transformer 层的权重矩阵的?以 Megatron 为例说明。

张量并行是将单层内的权重矩阵切分到多个 GPU 上并行计算,主要解决单卡显存无法容纳整个层参数的问题。

Megatron-LM 中的张量并行:

对于 Transformer 的自注意力和 FFN 层,权重矩阵按列或按行切分,配合通信操作。

MLP 块(两层全连接): 设输入为 X,权重为 A(升维)和 B(降维)。

  • A 沿列切分为 A1,A2,分别放在两张 GPU 上。输入 X 完全相同(由上一层输出广播而来)。各 GPU 计算 Yi=GeLU(XAi)。

  • B 沿行切分为 B1,B2。各 GPU 计算 Zi=YiBi。然后通过 AllReduce 将各 GPU 的 Zi 相加,得到最终输出。这里通信发生在激活值的 AllReduce。

image.png

通信分析:每次前向中,张量并行需要进行两次 AllReduce(分别在 FFN 后和注意力后)。通信量正比于批次大小 × 序列长度 × 隐藏维度。由于通信发生在层内,对带宽和延迟敏感,因此张量并行通常在节点内使用高速 NVLink 连接,不适合跨节点。

流水线并行(Pipeline Parallelism)如何划分模型层次?1F1B 调度策略是什么?

流水线并行将模型的不同层分配到不同 GPU 上。例如,一个 12 层的 Transformer,GPU0 负责第 1-4 层,GPU1 负责第 5-8 层,GPU2 负责第 9-12 层。

朴素流水线:将批次拆分为多个微批次(micro-batches)。GPU0 处理完第一个微批次后,将激活值传给 GPU1,然后立即处理第二个微批次,同时 GPU1 处理第一个微批次的剩余层。这样可以形成流水线,但存在“气泡”(bubble),即 GPU 空闲等待时间。

1F1B(one forward one backward)调度:这是 GPipe 的改进方案,用于减少气泡。其思想是:在稳态下,每个设备交替执行一个前向和一个反向(微批次),保持计算连续性。具体步骤:

  1. 预热阶段:每个设备按顺序接收微批次并执行前向计算,直到最后一个设备开始接收到第一个微批次。

  2. 稳态阶段:每个设备每完成一个微批次的前向,就紧接着执行一个微批次的反向,形成“前向-反向-前向-反向”的循环。

  3. 收尾阶段:不再有新的前向任务,各设备完成剩余的反向传播。

1F1B 通过将前向和反向交错执行,减少了设备的空闲时间,提高了硬件利用率。

什么是序列并行(Sequence Parallelism)?它主要解决什么问题?通常与哪种并行结合?

序列并行是将输入序列的长度维度切分到多个 GPU 上,以缓解长序列训练时激活值显存过大的问题。

在 Transformer 中,自注意力的计算复杂度与序列长度的平方成正比。当序列极长(如几千 token)时,注意力矩阵的显存占用成为瓶颈。序列并行通常在张量并行的基础上使用。例如 Megatron-LM 中,将序列长度切分,每张 GPU 只处理一部分序列 token。但在计算注意力时,需要获取完整的序列,因此需要使用 AllGather 通信来收集其他 GPU 的序列块;计算完注意力后,再通过 ReduceScatter 将结果分发回各 GPU。

序列并行通常与张量并行结合,因为两者都是切分模型层内计算,可以共用通信组,减少额外通信。它有效降低了单卡上注意力计算的峰值显存,使得训练更长序列成为可能。

ZeRO-1、ZeRO-2、ZeRO-3 分别对哪些状态进行分片?显存削减效果和通信开销各是怎样?

ZeRO(Zero Redundancy Optimizer)通过消除数据并行中的冗余存储来优化显存。

  • ZeRO-1:将优化器状态(如 Adam 的 momentum 和 variance)分片到各 GPU。每张卡只存储一部分参数的优化器状态,并负责更新这部分参数。通信:在反向传播后,需要对梯度进行 ReduceScatter,使每张卡拥有自己那部分参数的梯度,然后更新对应优化器状态和参数。之后,需要通过 AllGather 收集更新后的参数。显存削减:优化器状态部分几乎减少到 1/N(N 为 GPU 数)。通信量与标准数据并行相当(AllReduce vs ReduceScatter+AllGather)。

  • ZeRO-2:在 ZeRO-1 基础上,还将梯度分片。每张卡只存储自己那部分参数的梯度。反向计算时,每张卡只对负责的参数计算梯度。通信:类似,需要 ReduceScatter 将梯度分片到对应卡。显存节省更多(梯度也减少 1/N)。通信量略增。

  • ZeRO-3:进一步将模型参数也分片。每张卡只存储一层(或一部分)参数。前向需要某层时,通过 AllGather 从其他卡收集完整参数;计算后丢弃;反向同理。这几乎将显存占用均分给所有卡,允许训练超大模型。通信量大幅增加,因为每层前向和反向都需要 AllGather 参数和 ReduceScatter 梯度。

总的来说,ZeRO-1/2/3 逐步增加通信量以换取更高的显存效率。ZeRO-3 结合梯度检查点,可实现百亿甚至千亿参数模型的训练。

3D 并行:如何组合数据并行、张量并行、流水线并行来训练千亿甚至万亿参数模型?

3D 并行是同时使用三种并行策略:

  • 张量并行(TP):切分层内矩阵,用于单层过大时。通常在节点内使用高速 NVLink,TP size 一般 2-8。

  • 流水线并行(PP):将不同层分布到不同设备,用于模型过深时。PP 组跨节点,通信量较小。

  • 数据并行(DP):每个“模型副本”使用上述 TP+PP 组合,不同副本之间做数据并行。通常使用 ZeRO 优化数据并行中的冗余。

组合方式:将 GPU 划分为三维网格。例如,假设总 GPU 数为 G,设置 TP=4,PP=4,则每个模型副本占用 16 张 GPU,剩余的 GPU 用于数据并行(DP=G/16)。数据并行组内使用 ZeRO-2/3 来减少显存和通信。训练时,数据并行负责大批量样本的分布式处理,流水线并行负责层间流水线调度,张量并行处理层内大矩阵。

这种组合解决了单卡显存、计算和通信的极限,使得训练万亿参数模型成为可能。

大规模训练时,通信瓶颈常出现在哪里?如何用梯度压缩、异步通信来优化?

通信瓶颈:

  • 数据并行中的 AllReduce 梯度同步,当模型参数量巨大时,每步需传输梯度,跨节点网络带宽成为瓶颈。

  • 张量并行中的 AllReduce 激活值,频率高,对节点内带宽要求极高。

  • 流水线并行的点对点激活传输,跨节点时带宽需求也大。

优化方法:

  • 梯度压缩:通过量化(如 1-bit Adam)、稀疏化(只传输重要梯度)或低秩压缩(PowerSGD)来减少通信数据量。例如,将梯度从 FP32 压缩到 1-bit 符号,极大降低带宽。

  • 通信与计算重叠:利用 CUDA Stream,在反向计算的同时进行梯度通信,或在流水线中发送激活时同时计算。例如,DDP 的梯度 AllReduce 可以做到部分与反向传播重叠。

  • 异步通信:允许各 GPU 的参数更新存在一定延迟(异步 SGD),不再等待全局同步,但可能影响收敛。

  • 使用高速互联:节点内用 NVLink/NVSwitch,节点间用 InfiniBand 并启用 GPUDirect RDMA,避免 CPU 参与。

  • 结合 ZeRO:使用 ZeRO-1/2 时,ReduceScatter 和 AllGather 代替 AllReduce,通信量稍有增加但更利于分片。

如何选择并行策略的“最优”配置?比如张量并行度、流水线并行度和数据并行度如何根据硬件设定?

选择取决于硬件拓扑、模型规模和训练效率:

  • 节点内 GPU 间有高速 NVLink:优先使用张量并行,TP 大小通常等于节点内 GPU 数(如 8)。这样张量并行的 AllReduce 在节点内速度极快。

  • 模型层数很多(如 >64 层):使用流水线并行,PP 大小可跨节点,每个 PP 阶段放在一个节点内。PP 大小应使每个阶段的显存可容纳,同时气泡尽可能小(可使用 1F1B 调度)。

  • 剩余维度用于数据并行:DP 大小 = 总 GPU 数 / (TP × PP)。数据并行组内可使用 ZeRO 优化。

  • ZeRO 阶段选择:若显存允许,ZeRO-2 或 ZeRO-1 即可;大模型必须 ZeRO-3。

  • 序列并行:仅在序列极长且张量并行已经开启时启用,与 TP 结合。

最佳实践:先用少量 GPU 测试不同配置下的吞吐量(samples/sec),找到通信与计算的平衡点。Megatron-LM 提供了计算最优配置的工具。

模型分片后,不同 GPU 之间的负载不均衡怎么解决?

负载不均衡主要出现在流水线并行中,不同阶段的计算量可能不一致(比如嵌入层与中间层),导致某些设备成为瓶颈。

解决方法:

  • 自动平衡:框架(如 Megatron-LM)将层尽可能均匀分配到流水线阶段,同时考虑参数量和计算量。

  • 微批次调度优化:调整流水线气泡,如使用 interleaved 1F1B,每个设备处理多个流水线阶段(交错),进一步平衡负载。

  • 使用动态批大小或序列长度:但可能复杂。

  • 考虑给计算量小的阶段增加虚拟层:但实际上很少。

对于张量并行和数据并行,只要切分均匀,负载自然均衡。ZeRO 分片也是均衡的。

在大规模集群训练中,如何实现高效的检查点保存和恢复?有什么最佳实践?

大规模训练中,检查点保存巨大(数十 GB 至 TB 级别),需要高效且容错。

  • 分布式原子保存:每个 GPU 仅保存自己分片的那部分模型参数和优化器状态,而不是完整模型,减少单卡 I/O。使用 DeepSpeed 或 Megatron 的分布式 checkpoint 格式。

  • 异步写入:计算与磁盘写入重叠,利用独立的 CUDA Stream 或 CPU 线程。

  • 快速存储:使用高速共享存储(如 NVMe SSD 集群、Lustre、GPFS),或本地 NVMe + 后台上传。

  • 增量保存:仅保存自上次 checkpoint 之后的参数变化(难度较大,现仍多全量保存)。

  • 恢复流程:确保所有分片文件就绪,需要时每卡读取自己分片,并验证分片一致性。

  • Warm 重启:恢复时不仅要加载模型参数,还要恢复优化器状态、学习率调度器、数据迭代器状态,以精确复现训练。

  • 监控和自动故障恢复:结合集群调度器,在节点故障时自动从最近检查点恢复,利用弹性训练动态调整 GPU 数量。

多机多卡训练时,NCCL 通信库的作用是什么?有哪些常见问题?

NCCL(NVIDIA Collective Communications Library)是专为多 GPU 和多节点环境设计的高性能通信原语库,提供了 AllReduce、AllGather、ReduceScatter、Broadcast 等集合操作,并针对 NVLink、PCIe、InfiniBand 等硬件进行了拓扑感知优化。它是 PyTorch DDP、DeepSpeed、Megatron 等框架的底层通信引擎。

常见问题:

  • 通信死锁:由于代码逻辑错误导致各 GPU 执行不同的通信操作,造成相互等待。需确保所有 rank 以相同顺序调用集合操作。

  • 性能瓶颈:未启用 GPUDirect RDMA,数据需要经过 CPU 中转,延迟增大。需要配置环境变量 NCCL_NET_GDR_LEVEL 等。

  • 版本不兼容:NCCL 版本与驱动、CUDA 版本不匹配可能导致崩溃或性能下降。

  • 环拓扑问题:在 InfiniBand 网络上,错误的路由设置可能使 Ring AllReduce 效率低下。

  • 超时:大规模集群中,慢节点或高通信负载可能导致 NCCL 操作超时,需要适当调整超时设置。

解释一下 ZeRO-Offload 如何利用 CPU 内存和计算来进一步降低 GPU 显存需求。

ZeRO-Offload 将优化器状态(甚至模型参数)卸载到 CPU 内存,利用 CPU 的算力来执行优化器更新。具体地:

  • 优化器状态(如 Adam 的 momentum 和 variance)体积庞大,存储在 CPU 内存中。

  • 在反向传播结束后,GPU 将梯度发送给 CPU,CPU 执行 Adam 更新步骤,得到更新后的 FP32 参数。然后 CPU 将更新后的参数传回 GPU,GPU 再将其转换为 FP16 用于下一轮。

  • 为了隐藏通信延迟,利用双缓冲和异步传输:在 CPU 计算当前层的更新时,GPU 可以预取下一层的梯度,或者进行下一轮的前向。

  • 还可以将 FP16 模型参数也部分放在 CPU,在需要时通过 PCIe 传输到 GPU,实现“零冗余”参数存储。

代价是训练速度下降,因为 PCIe 带宽远小于 GPU 显存带宽。典型性能损失在 30%~50%,但可将单卡可训练模型规模扩大数倍。

对于 MoE 模型,数据并行和专家并行如何结合?

MoE(Mixture of Experts)模型中,FFN 层被多个专家替换,每个 token 通过路由器(gating)被分配到少数几个专家。

  • 专家并行(Expert Parallelism):将不同的专家分布在不同的 GPU 上。输入 token 根据路由结果,通过 AllToAll 通信发送到对应的专家所在 GPU 进行计算,输出后再通过 AllToAll 返回。这本质上是将专家维度进行切分。

  • 数据并行:处理非专家部分(如注意力层)以及不同输入样本。通常与专家并行组合。

组合方式:模型副本之间采用数据并行,每个副本内部对专家使用专家并行。例如,8 张 GPU,分成 2 个数据并行组(每组 4 张),每组内 4 张 GPU 进行专家并行,即每个专家被放置在一张卡上,全组 4 张卡共同拥有所有专家。训练时,数据并行组独立处理不同批次,组内专家并行进行 AllToAll 通信。

这种结合有效扩展了 MoE 的训练规模,平衡了通信与计算。

你从 0 搭建一个千卡训练平台,技术选型和调试流程会是怎样的?

技术选型:

  • 硬件:NVIDIA A100/H100 GPU,节点内 NVSwitch 全互联,节点间 InfiniBand HDR/NDR 网络。

  • 操作系统:Ubuntu 或 Rocky Linux,安装 MLNX_OFED 驱动、NVIDIA 驱动、CUDA、cuDNN、NCCL。

  • 容器化:使用 Docker/Enroot + Slurm/Kubernetes 进行作业调度。

  • 训练框架:PyTorch + DeepSpeed(ZeRO)或 Megatron-LM(3D 并行)。对于千卡规模,通常使用 Megatron 进行模型并行,结合 DeepSpeed ZeRO 优化数据并行。

  • 存储:高性能共享文件系统(Lustre/GPFS)用于数据和检查点。

  • 监控:Prometheus + Grafana 监控 GPU 利用率、温度、功耗、网络流量;NVIDIA DCGM 暴露 GPU 指标。

调试流程:

  1. 单卡验证:确保模型、数据加载、训练循环在单卡上正常运行。

  2. 单机多卡(8 卡):使用 DDP 验证多卡训练正常,无死锁,性能线性。

  3. 小规模集群:选取 2-4 个节点,配置张量并行+流水线并行+数据并行,调通分布式 checkpoint 和弹性训练。

  4. 全规模启动:逐步扩展到所有节点。使用较小的模型和数据集进行“冒烟测试”,检测通信瓶颈、慢节点。利用 nccl-tests 测试带宽和延迟。

  5. 性能剖析:使用 PyTorch Profiler、Nsight Systems 定位瓶颈,调整并行策略、微批次大小、梯度累积步数、通信与计算重叠。

  6. 稳定性测试:连续运行数小时,监控 loss 变化、GPU 内存、通信超时等,确保无随机故障。

  7. 正式训练:加载完整数据,开始正式预训练,配合检查点保存和自动恢复机制。

整个过程需要深入理解硬件拓扑和框架实现,并不断优化通信和计算效率。