第 11 章 分布式训练框架与系统
第 11 章 分布式训练框架与系统
学习目标
- 掌握主流分布式框架的分工:PyTorch DDP/FSDP、DeepSpeed、Megatron-LM、Horovod/Ray Train/Colossal-AI;
- 理解 NCCL 与通信拓扑(NVLink、InfiniBand、RoCE);
- 理解 Checkpoint、断点续训、弹性训练的设计;
- 理解千卡/万卡训练的核心度量:MFU、吞吐、扩展效率与成本。
11.1 框架格局:谁负责什么
分布式训练框架分三层:通信库(NCCL)→ 并行封装(DDP/FSDP/DeepSpeed/Megatron)→ 调度(第 12 章)。选框架的本质是选「并行策略的封装程度」:
| 框架 | 策略 | 定位 | 适用 |
|---|---|---|---|
| PyTorch DDP | 数据并行(朴素) | 官方,简单 | 模型单卡可放,多卡加速 |
| PyTorch FSDP | ZeRO-3 分片 | 官方,生态好 | 7B-70B 微调/预训练,社区主流 |
| DeepSpeed | ZeRO 全系 + 混合并行 + Offload | 功能最全 | 大规模微调、MoE、异构卸载 |
| Megatron-LM | TP+PP+SP+EP 原生组合 | 大厂级预训练 | 100B+ 预训练、性能极致 |
| Horovod | 跨框架数据并行 | 历史主流 | 老代码库、多框架混用 |
| Ray Train | 调度 + 训练封装 | 弹性/编排 | 超参搜索、多任务调度 |
| Colossal-AI | 一键并行策略 | 易用性 | 快速原型、中小团队 |
现实格局:预训练大模型 = Megatron 系(或 NVIDIA 的 Megatron-DeepSpeed 组合);微调/中小模型 = PyTorch FSDP 与 DeepSpeed ZeRO 为主;编排层 = Ray/KubeRay(第 12 章)。
# FSDP:一句话开启 ZeRO-3 风格分片(PyTorch 官方)
from torch.distributed.fsdp import FullyShardedDataParallel as FSDP
from torch.distributed.fsdp.wrap import transformer_auto_wrap_policy
model = FSDP(
model,
auto_wrap_policy=transformer_auto_wrap_policy(TransformerBlock),
sharding_strategy="FULL_SHARD", # ZeRO-3 等价
)11.2 DeepSpeed:ZeRO 之外的武器库
DeepSpeed(微软)不只是 ZeRO。生产里常用它的组合能力:
- ZeRO-Offload / NVMe Offload:把优化器状态或参数卸载到 CPU 内存/NVMe,让单机多卡跑 70B 级微调(第 8 章「最后手段」的实现);
- 混合并行:ZeRO + Megatron 风格 TP/PP 的组合(Megatron-DeepSpeed 是 100B+ 预训练的主流组合之一);
- MoE 优化:DeepSpeed-MoE 内置专家并行与通信优化;
- 通信压缩:梯度压缩(1-bit Adam 等)降低通信量。
DeepSpeed 的配置文件驱动风格:
# ds_config.json 片段
{
"train_batch_size": 1024,
"gradient_accumulation_steps": 16,
"zero_optimization": {
"stage": 3,
"offload_optimizer": {"device": "cpu"}
},
"bf16": {"enabled": true}
}11.3 Megatron-LM:TP/PP/SP 的原生实现
Megatron-LM(NVIDIA)的价值不在「另一个 ZeRO」,而在把张量并行、流水并行、序列并行做成同一套代码的一等公民:
- TP 实现:列并行(Column Parallel,切 的列)+ 行并行(Row Parallel),配合两个 all-reduce(f 与 g 的对称设计,省一次通信);
- PP 实现:1F1B 交错调度(第 10 章),配合梯度累积;
- SP:LayerNorm/Dropout 的序列切分,省激活显存;
- MoE 支持:EP 作为一等并行维度。
Megatron 的核心主张:TP/PP/SP 在 Megatron 里是「配置即用」而不是「手写拼装」,这让 100B-1T 模型的训练从「系统研究问题」变成「配置问题」。这也是它与 PyTorch 生态的裂缝所在——Megatron 维护自己的一套张量库(Parallel Transformer),社区用 FSDP 更顺。
11.4 通信拓扑:NVLink、InfiniBand、RoCE
集合通信跑在什么网络上,决定分布式训练的天花板:
| 层 | 技术 | 带宽量级 | 延迟 | 用途 |
|---|---|---|---|---|
| 机内 | NVLink(+NVSwitch) | 900GB/s 级 | 微秒级 | TP、同机梯度同步 |
| 机间 | InfiniBand(IB) | 400Gbps-800Gbps | 微秒级 | DP/PP/EP 跨机 |
| 机间 | RoCEv2 | 同 IB 量级 | 略高 | 以太网上的 RDMA,成本低 |
| 机间 | 普通以太网 | 100Gbps 级 | 高 | 仅控制面/日志 |
三个工程要点:
- NCCL 是事实通信标准:NVIDIA 的 NCCL 库为各拓扑做了算法选择(ring/tree/ragged),把「用什么算法跑 AllReduce」这件事从手写中解放。AMD 用 RCCL,华为昇腾有自己的 HCCL。
- 拓扑感知(Topology-Aware):NCCL 会感知 GPU 在机架/交换机的位置选择最优通信路径;坏一条 IB 链路,整个集群训练都可能被拖慢——网络是集群健康的第一敏感点。
- 三网分离(第 12 章展开):训练数据、集合通信、管理监控走不同的网络平面,避免互相干扰。
11.5 Checkpoint、断点续训、弹性训练
11.5.1 Checkpoint 的设计
大模型训练跑数周,任何一次故障都要求能回到最近的良好状态。Checkpoint 要保存的东西(第 4 章显存账的镜像):模型权重(BF16)+ 优化器状态(FP32 矩)+ 学习率调度器位置 + RNG 状态(数据采样随机性,保证恢复后可复现)。
- 全量 Checkpoint:每 N 步存一份,7B 约 15GB,70B 约 150GB——写盘本身就占训练时间(存储带宽是瓶颈);
- 异步 Checkpoint:权重先异步刷盘,不阻塞训练步进;
- 分片 Checkpoint:FSDP/ZeRO 各卡只存自己的分片,恢复时再聚合——省时省空间,但恢复要等聚合。
11.5.2 断点续训与故障恢复
断点续训(Resume)不是「从 checkpoint 加载」这么简单,三个雷区:
- 数据顺序:恢复后数据迭代器要从正确位置继续(否则重复学同一批数据);
- 优化器状态:优化器矩必须一起恢复(只恢复权重 = 重新预热);
- RNG 状态:dropout 与数据采样的随机序列要恢复,否则「复现」变味。
弹性训练(Elastic Training):节点动态加入/退出(云上抢占式实例、故障替换)时不中断训练——把「重启整个任务」变成「重算掉队节点」。TorchElastic(PyTorch)、Ray 都提供弹性能力。工程权衡:弹性引入的复杂度 vs 云上成本收益,固定集群往往不需要,抢占式集群几乎必须。
11.6 千卡/万卡训练:MFU、吞吐、扩展效率、成本
11.6.1 MFU 与吞吐
MFU(Model FLOPs Utilization):实际吞吐 ÷ 硬件理论峰值 FLOPs,即有效吞吐 与硬件峰值 之比:
- 单卡 MFU 的天花板来自算子效率与 IO(H100 FP16 理论 ~1000 TFLOPS,实际好的 kernel 跑到 60-70% 就很好);
- 分布式后 MFU 下降是必然:通信、气泡、负载不均都在吃掉算力。千卡集群 MFU 40% 是优秀,20% 是常态。
配套指标还有 Model TFLOPs/GPU/s(每卡每秒有效模型计算量,跨规模可比)与 吞吐(token/s 或 samples/s)。
11.6.2 扩展效率(Scaling Efficiency)
其中 为 卡时的有效吞吐, 为单卡有效吞吐。
8 卡 0.9 很常见,1024 卡 0.5-0.7 已经不错。扩展效率随规模下降的三个原因:通信占比上升(第 10 章)、负载不均(MoE 的 hot expert)、掉队节点。
%%{init: {"themeVariables": {"primaryTextColor": "#000000", "textColor": "#000000", "labelColor": "#000000", "nodeTextColor": "#000000", "labelTextColor": "#000000", "scaleLabelColor": "#000000"}}}%%
flowchart LR
A[16 卡
效率 0.85] --> B[128 卡
效率 0.70]
B --> C[1024 卡
效率 0.55]
C --> D[4096 卡
效率 0.40]
D --> E[10000+ 卡
效率 0.30]
style E fill:#fde8e811.6.3 成本账与「该不该买卡」
训练成本 = 卡数 × 时长 × 单价。用第 4 章的 FLOPs 账可以倒推「该不该用更多卡」:
- 加速比 = 卡数 × 扩展效率。加卡到扩展效率跌破某个阈值后,加卡反而更贵;
- 云上还要考虑:预留实例 vs 抢占实例(抢占便宜 60-80%,但要弹性训练兜底,第 12 章);
- 训练 vs 推理的成本结构完全不同(第 17 章):训练是一次性成本,推理是持续性成本——多数产品最终花在推理上的钱远多于训练。
本章要点回顾
- 框架格局:预训练用 Megatron,微调用 FSDP/DeepSpeed,编排用 Ray;NCCL 是通信事实标准。
- DeepSpeed 的价值在组合能力:ZeRO 分片 + Offload + MoE + 通信压缩。
- Megatron 把 TP/PP/SP 做成配置,是 100B+ 预训练的系统答案。
- 通信拓扑决定天花板:NVLink 管机内(TP)、IB/RoCE 管机间(DP/PP/EP);三网分离防干扰。
- Checkpoint 要保存权重+优化器+调度器+RNG;断点续训的雷区是数据顺序与优化器状态;弹性训练服务抢占式集群。
- MFU 衡量算力利用(千卡 40% 已优秀);扩展效率随规模下降是物理规律,不是 bug。
- 成本账的关键认知:加卡收益递减、抢占实例要弹性兜底、推理是持续性大头。
习题
- 你的 7B 微调任务单机 8 卡 FSDP 只有 30% 扩展效率。列出 5 个可能瓶颈并给出排查顺序。
- 计算:MFU 40% 的 1024 卡 H100 集群(每卡 ~1e15 FP16 FLOPs/s),训练 70B 模型 1.5 万亿 token 需要多少天?写出 FLOPs 账。
- 断点续训后 loss 比断点前高 0.1 且不回降,最可能是什么没恢复?如何验证?
- 为什么 MoE 训练的扩展效率比稠密模型更难保持?从 hot expert 与 all-to-all 角度分析。
- 设计一个抢占式实例集群的训练容错方案:检查点频率、失败检测、恢复流程分别怎么定?
延伸阅读
- Rajbhandari et al., ZeRO: Memory Optimizations Toward Training Trillion Parameter Models, 2020
- Shoeybi et al., Megatron-LM: Training Multi-Billion Parameter Language Models Using Model Parallelism, 2020
- Narayanan et al., Efficient Large-Scale Language Model Training on GPU Clusters Using Megatron-LM, 2021
- Zhao et al., PyTorch FSDP: Experiences on Scaling Fully Sharded Data Parallel, 2023
- Rasley et al., DeepSpeed: System Optimizations Enable Training Deep Learning Models with Over 100 Billion Parameters, 2020
- NVIDIA NCCL Documentation(集合通信算法与拓扑)
- 各模型技术报告中的 MFU/扩展效率数据:Llama 3(Meta)、Qwen(Alibaba)、DeepSeek-V3
下一章预告
第 12 章从「训练本身」跳到「训练跑在什么上」:Kubernetes/Slurm/Volcano 调度、GPU 资源管理、多租户配额,以及实验管理、日志、指标、追踪——训练基础设施的全景。