Back to blog

Training LLMs with Fault Tolerant HSDP on 100,000 GPUs

Meta 提出的容错混合分片数据并行 FT-HSDP:以数据并行副本为容错单元,通过 CPU-GPU 混合 FTAR 全归约与非阻塞追赶协议实现异步恢复,在 10 万 GPU 规模将有效训练时间从 44% 提升至 80%。

Training LLMs with Fault Tolerant HSDP on 100,000 GPUs:面向 10 万 GPU 的容错训练范式

一、论文概述

项目内容
标题Training LLMs with Fault Tolerant HSDP on 100,000 GPUs
作者Omkar Salpekar, Rohan Varma, Kenny Yu, Vladimir Ivanov, Yang Wang, … , Chunqiang Tang, Mathew Oldham(共 24 位作者)
机构Meta
论文https://arxiv.org/abs/2602.00277
代码未开源
发布2026-01-30 (v1)
规模生产环境 98K H100 GPU,跨 4 座数据中心楼宇

一句话概括

当训练集群规模达到 O(100K) GPU 时,同步训练的故障恢复开销高到无法接受——每 18 分钟就发生一次故障,而每次同步恢复要停机约 10 分钟,有效训练时间仅 44%。本文提出 FT-HSDP(Fault Tolerant Hybrid-Shared Data Parallelism,容错混合分片数据并行),以数据并行副本(replica)为容错单元,配合 CPU-GPU 混合的容错全归约协议(FTAR) 与 非阻塞追赶(non-blocking catch-up)协议,实现异步恢复:单个副本故障时其余副本继续训练,将停机时间从 10 分钟降到 3 分钟,有效训练时间提升至 80%,且不显著损害模型精度。


二、核心思想

问题定义

大规模同步训练(fully synchronous training)要求所有 GPU 同时健康才能推进训练。一旦某个 GPU 或服务器故障,所有节点都必须从上一个检查点重启。此外,NCCL 等集合通信库要求初始化时已知全部成员,不支持动态重配置,因此重启还需重新初始化所有通信连接,在超大规模下造成长时间停机。

规模从 32K 扩展到 100K GPU 时,同步训练面临两个叠加的挑战:

  1. 故障率随规模上升:据 Meta 生产数据估算,100K GPU 下平均每 18 分钟发生一次故障。
  2. 恢复停机时间随规模上升:即使投入大量工程优化,100K GPU 下的同步恢复仍需约 10 分钟。

于是每 18 分钟里,10 分钟花在故障切换与重启上,仅剩 8 分钟用于有效训练:

有效训练时间=18−1018=818≈44%\text{有效训练时间} = \frac{18 - 10}{18} = \frac{8}{18} \approx 44\%

这是不可接受的。

解决方案概述

FT-HSDP 沿用 HSDP 思路,将全部 GPU 划分为多个副本(replica),每个副本由数千 GPU 组成(内部混合使用数据/张量/流水线/专家/上下文并行),负责训练一部分输入数据。不同副本之间以数据并行方式周期性交换梯度。

副本化带来两个容错机会:

  1. 缩小恢复范围:故障发生时,只需重建包含故障节点的那一个副本,而非整个集群,恢复规模更小、恢复时间更短。
  2. 异步恢复:故障副本恢复期间,其余健康副本可以继续训练,而不是全体停机等待。

尽管副本化容错的想法在此前工作(如 Horovod、ReCycle)中已有探索,但本文回答了两个关键工程问题:

  • Q1:如何在 10 万 GPU 规模真正落地这一思路?(尤其是 NCCL 无法动态增删成员的问题)
  • Q2:异步恢复带来的”每步训练数据量不固定”是否会损害精度与收敛?

三、背景:集群网络与硬件可靠性

3.1 多楼宇网络拓扑

集群网络架构-楼内

10 万 GPU 集群必然跨越多座数据中心(DC)楼宇。Meta 设计了可将数十万 GPU 整合进单一高性能 RoCE fabric 的多楼宇网络:

  • 楼内(Figure 1a):采用 3 层 Clos 架构。
    • RTSW(Rack Training Switch):连接机架内 GPU
    • CTSW(Cluster Training Switch):连接一个 AI Zone 内所有机架
    • ATSW(Aggregator Training Switch):跨 AI Zone 连接各 CTSW,将 RoCE 网络扩展到单个 AI Zone 之外
    • 跨 AI Zone 超额订阅比(over-subscription ratio)为 1:2.8

集群网络架构-跨楼

  • 跨楼(Figure 1b):不同 DC 楼宇的 ATSW 层之间采用全连接网格(fully connected mesh),跨 DC 流量的超额订阅比同样为 1:2.8,可增量扩展到同一 RoCE fabric 内的数十万 GPU。

关键延迟特性:随网络跳数增加,GPU-to-GPU 通信延迟显著增大。相对于同机架 GPU(最低延迟),跨机架同 AI Zone、跨 AI Zone、跨 DC 楼宇分别为 7×、15×、30× 更高延迟。

拓扑对设计的启示:这一延迟特性天然引导出副本化数据并行设计——将一个副本的全部 GPU 放在同一个 AI Zone 或 DC 内(延迟敏感的副本内集合通信不必走跨 DC 链路);将不同副本放在不同 DC(跨副本的数据并行集合通信对高延迟更鲁棒)。

3.2 硬件可靠性数据(32K GPU 生产统计)

Table 1 汇总了一个约 32K GPU 常态训练任务的中断分类(有效训练时间 95%–97%,平均每 1000 台服务器每天 2.3 次中断):

故障类别次数占比
GPU HBM3 内存15522.9%
PCIe 设备12218.0%
NCCL 看门狗超时619.0%
GPU 计算故障507.4%
软件 Bug487.1%
主机维护426.2%
内核故障395.8%
系统重启385.6%
数值/静默数据损坏(SDC)375.5%
网络交换机/线缆365.3%
SSD304.4%
其他(SRAM/未知/内存/散热/传感器)20~3%

关键观察:

  • 78% 的中断源于硬件故障。
  • 过去排名第一的 “GPU 计算故障” 经前期治理已降至第四(7.4%);HBM 问题升为首位(22.9%)。
  • PCIe/SSD 问题激增:SSD 因 NVMe 固件 fatal bug 进入只读模式并报 I/O 错误。
  • 数值问题/SDC(静默数据损坏)最难处理:37 次中断溯源到 7 台主机,根因排查耗时数小时到数天。一个 SDC 案例中,某层某专家的 32 个连续参数(2T 参数中)突然飙升超过 1e7。

确定性训练(Deterministic Training) 用于定位 SDC:保证 (1) 任意训练迭代的数据跨运行一致;(2) 所有算子对相同输入给出确定性输出(通过禁用 flash attention 优化、强制 cuDNN 确定性算法、关闭 cuDNN auto-tuner 实现)。这样固定迭代数后可生成等价检查点,通过跨多次运行比对相同 rank 的检查点即可发现异常主机。

3.3 快速根因分析

  • 工具每 5–20 秒从所有 GPU 拉取追踪数据做深度分析。
  • 挂起(hang)诊断:先构建集合通信的 wait-for 图,叶子节点即整个任务阻塞的集合;再对每个集合构建 GPU 级 wait-for 图,判断 GPU 是”未加入集合”还是”加入但无进展”。
  • 性能下降诊断:聚合各 rank 的集合完成时间,完成时间显著偏低的 GPU 通常是拖累者(因为它最后才加入集合)。
  • 单次完整分析在 32K+ GPU 上 <5 秒完成,遥测数据驻留主机 DRAM(非 GPU HBM),每 GPU <1MB,异步采集不影响训练性能。
  • 注入约 1500 次故障验证,97.8% 的情况正确识别出故障服务器,其余 2.2% 也被纳入小候选集。

四、挑战:为何同步恢复在 100K 规模行不通

恢复时间分解

在 32K GPU 下同步训练尚可接受(约损失 3–5% GPU 时间)。但扩到 100K GPU 时故障率与恢复成本双双上升。Figure 2 分解了恢复过程各阶段耗时:

恢复阶段说明随规模变化
寻找替换 GPU(Allocation)naive “cold-N”:停整个任务、归还所有 GPU 给私有云、请求重启,可达 5 分钟。优化 “warm-N”:预留少量备用 GPU 替换故障 GPU,大幅降低开销cold 版本极高
NCCL 初始化重启训练需重新初始化 NCCL,建立所有连接的时间随 GPU 数增长:16K GPU 需 17 秒 → 98K GPU 约 200 秒(已基于定制 NCCL 优化)随规模显著增长
其他加载开销启动 PyTorch、拉取并加载检查点、健康检查等不随规模显著增长
”首步”开销(First-step)首步需额外任务:创建 checkpointer、初始化 dataloader(缓存后续 N 批)、JIT 编译等。首步可达数分钟,而常态步约 20 秒cold 版本高度不确定

结论:尽管投入大量优化,仍无法在更大规模保持恢复时间稳定。理想情况下恢复时间应随规模下降以抵消更高的故障率,但实际恰恰相反。100K GPU 下每 18 分钟一次故障、每次停机 10 分钟 → 有效训练时间仅 44%。这自然驱动团队转向异步恢复:异步恢复可以隐藏 allocation、NCCL 初始化和其他加载步骤的延迟(首步的额外延迟可能无法完全隐藏,但可将部分开销挪入恢复过程)。


五、FT-HSDP 设计

5.1 整体架构

FT-HSDP 架构

FT-HSDP 创建多个副本,每个副本处理一部分训练数据。每个副本内每个 GPU 有唯一 rank 号。训练完一个固定量数据(batch)后,不同副本中相同 rank 的 GPU 之间交换梯度(跨 DC 执行),最后各 rank 的优化器应用梯度更新权重。

工作流程(以”步 step”为单位):

  1. 每个副本对自己的 batch 做前向+反向计算与通信;
  2. 相同 rank 的 GPU 交换梯度求和(FTAR,与反向计算重叠以提速);
  3. 各 rank 优化器应用梯度更新权重;
  4. 可选:步末做检查点,将自身状态写入持久存储。
┌──────────────────────────────────────────────────────────┐
│                      FT-HSDP (跨 DC)                     │
│   Replica 0     Replica 1     Replica 2     Replica 3      │
│   (DC-A)        (DC-B)        (DC-C)        (DC-D)         │
│  ┌────────┐   ┌────────┐   ┌────────┐   ┌────────┐        │
│  │Rank 0  │←─→│Rank 0  │←─→│Rank 0  │←─→│Rank 0  │  FTAR  │
│  │Rank 1  │←─→│Rank 1  │←─→│Rank 1  │←─→│Rank 1  │ 相同   │
│  │  ...   │   │  ...   │   │  ...   │   │  ...   │ rank   │
│  │Rank k  │←─→│Rank k  │←─→│Rank k  │←─→│Rank k  │ 交换梯度│
│  └────────┘   └────────┘   └────────┘   └────────┘        │
│  副本内 NCCL   副本内 NCCL   副本内 NCCL   副本内 NCCL       │
│  (DP+TP+PP+CP)                                             │
└──────────────────────────────────────────────────────────┘

5.2 Rank 0(Leader)设计

Rank 0 Leader 设计

Rank 0 是一个副本的 leader。除了和其他 rank 一样有多条副本内计算/通信 stream、一条跨副本梯度交换 stream、主线程执行优化器步(OPT)外,还承担额外控制逻辑:

  1. 一个 CPU 线程通过共识服务(类似 Chubby 与 Delos)与其他副本的 Rank 0 协调,确定哪些副本健康;
  2. Rank 0 的主 GPU 线程询问同副本其他 rank 是否成功完成梯度交换,若是则可推进到优化器步。

注意 FT-HSDP 将 FTAR 与反向计算重叠以提升训练速度。

5.3 三个核心机制

(1)故障检测

当前实现依赖**超时(timeout)**检测故障 GPU。一步平均约 20 秒,故超时阈值经验设为 60 秒。未来计划与快速根因分析组件(§3.3)集成以进一步缩短停机。

(2)故障后的一致性保证:副本内 2PC + 跨副本无需一致

理想情况下所有 GPU 应处于一致状态。同步系统靠”故障即全体从最新检查点重启”来保证。但 FT-HSDP 中某 GPU 故障时,其余副本可能处于不一致状态(例如 Replica 0 的 Rank 0 故障前已与 Replica 1 交换梯度,但尚未与 Replica 2、3 交换)。

核心洞察:副本内一致性是必要的,但跨副本一致性不必要——因为 FT-HSDP 支持异步恢复。若一个副本完成了某步而另一个没完成,让完成的副本继续训练,未完成的副本稍后像恢复副本一样重新加入。

具体做法(类似两阶段提交 2PC):故障检测后,副本内 Rank 0 询问每个其他 rank 是否已完成梯度交换:

  • 全部回答”是” → Rank 0 令所有 rank 应用梯度更新权重,推进下一步;
  • 否则 → Rank 0 令所有 rank 重试该步。此时因为权重尚未更新,各 rank 只需丢弃梯度而无需从他处拉取模型权重。

由于每个副本独立决策,可能出现不同状态(如 Replica 1 推进到 step 100,而 Replica 2、3 重训 step 99)。落后的副本随后触发恢复协议。

(3)异步恢复

FT-HSDP 在其余副本继续训练的同时恢复某个副本。恢复分两种情况:

  1. 含故障 GPU 的副本:找替换 GPU、重建副本、重启所有 GPU,运行非阻塞追赶协议加入其他副本;
  2. 因故障后不一致而落后的副本:执行相同的非阻塞追赶协议加入。

5.4 高效容错全归约 FTAR

FT-HSDP 副本内通信仍用 NCCL(性能最优),但跨副本梯度交换引入自研的 FTAR(Fault Tolerant All Reduce) 协议,因为 NCCL 无法满足以下目标:

  • 网络感知架构:FTAR 跨 DC 执行,不能涉及 NVLink 域通信;跨楼带宽超额订阅,须考虑拥塞控制。
  • 可重建的通信组:故障时每个 FTAR 组须在健康 rank 上重建(清理陈旧传输资源、与新 peer 组重建连接)。
  • 易于故障管理:副本故障时健康副本的 rank 须保持存活并在内存中保留训练状态。FTAR 须能区分可恢复错误(如对端死亡导致的网络远程错误 → 上报并移除死 peer)与致命错误(如坏 NIC → 指示 kill 该 rank)。
  • 对并发计算 kernel 最小资源竞争:FTAR 与反向计算重叠,须最小化 GPU 资源竞争。

FTAR 协议细节

CPU-GPU 混合设计(核心创新)

根本问题:NCCL 由 GPU 驱动通信,但 GPU 作为主要面向大规模并行的设备,缺乏 CPU 实现复杂控制逻辑的能力(区分错误类型、动态增删连接、拥塞控制等)。此前的原型和其他工作(Horovod、RobustLLM)需要重启所有节点来重建连接,耗时数分钟,几乎抵消异步恢复的收益。

关键观察:用 HSDP 副本作为容错单元简化了重配置——因为跨副本只有一个集合通信(all reduce)。

FTAR 采用 CPU 驱动控制平面 + GPU 执行数据平面 的混合设计(Figure 5):

平面承担方职责
控制平面CPU初始化 RDMA 连接(创建发送/接收缓冲区);运行 reconfig 确定哪些 rank 加入 FTAR 组;维护连接(可自由销毁/新建);通过共识服务确定 peer;限制在途数据量做拥塞控制;按错误类型决定处理方式
数据平面GPUGPU stream 将数据拷入发送缓冲区并通知 CPU asyncThread;CPU asyncThread 通过 RDMA 发送数据并等待其他 rank 数据;数据到达后通知 GPU stream 执行 reduce。所有数据传输由 GPU 执行以最大化性能

CPU-GPU 交互流程:

1. CPU 主线程初始化 RDMA 连接(send/recv buffer)
2. CPU 运行 reconfig 确定 FTAR 组成员
3. 启动 GPU stream 和 CPU asyncThread
4. GPU stream: 拷贝数据到 sendBuf → 通知 CPU
5. CPU asyncThread: 通过 RDMA 发送数据到远程 rank
6. CPU asyncThread: 等待其他 rank 数据到达
7. CPU 通知 GPU stream 数据已到达
8. GPU stream: 对 recvBuf 执行 reduce
9. 可多轮迭代,每轮交换一部分梯度(拥塞控制)

Ring 算法与流水线

FTAR 参与的 GPU 大多分布在跨 AI Zone、跨楼(高超额订阅、有限对分带宽),须避免网络拥塞。因此选用 Ring 算法(每个 GPU 只与环上两个邻居收发,最小化组内并发流量),针对 200MB–500MB 消息、最多 16 rank 优化(偏好带宽最优算法)。

Ring 算法分两阶段,各 N-1 步(N = 环中 GPU 数 = 副本数):

  • ReduceScatter 阶段:副本 i 的 GPU 先发送第 i 份数据给右邻居;每个 GPU 收到数据后做 GPU 内 reduction 再转发。最后每个 GPU 拥有一份数据的最终结果。
  • AllGather 阶段:每个 GPU 将其最终结果份转发给右邻居并逐跳传递,最终所有 GPU 获得全部份。

固定大小 chunk 流水线:初始化时在每个 GPU 预分配内部 sendBuf 和 recvBuf,大小为 固定 chunk 大小 (S) × chunk 数 (C),并预注册到网络用于 RDMA。对 N rank 的 FTAR 组,数据切成多个分区,每个分区至多 S×C×NS \times C \times N 字节;分区内跑一轮 2N-2 步 Ring 并用 C 个 chunk 流水线。两个好处:(1) 控制每两个 peer 间并发数据量至多 S×CS \times C 字节;(2) 允许独立调优 copy/reduce kernel(thread block 数)与网络传输(Queue Pair 数等)。

Kernel 优化(低占用高带宽)

  • ReduceScatter 步 0~N-2 结果只存 sendBuf;第 N-1 步存到本地数据缓冲(返回用户)和 sendBuf。合并这两步减少一轮 CPU-kernel 同步,且避免后续 HBM load(结果仍在寄存器中)。
  • 算法为每个 AllReduce 启动一个 CUDA kernel,在 host-pinned 内存的 flag 上忙轮询。为减少 GPU 资源浪费,kernel 在极少数 SM 上启动(H100 有 132 SM,NCCL AllReduce 仅用 4 SM)。为在低占用下实现高内存总线利用率,用**指令级并行(ILP)**每线程发多条 load/store/reduce 指令,并通过多种手段减少寄存器占用(constexpr 编译期计算、减小 block size 等)。
  • 最终配置:FTAR 用 2 个 thread block(2 SM)、每 block 512 线程、8MB chunk,把 reduction 完全隐藏在并发网络传输之下——用的 SM 比原生 NCCL AllReduce(4 SM)还少。因为 AllReduce 本就是网络受限,多用 SM 反而严重降低其他计算 kernel 性能。

超时与错误处理

GPU stream 等待 peer 数据/信号时用带超时的 while 循环检查 flag,确保 GPU 不会挂起。错误分两类:

错误类型例子处理
可恢复错误超时、ncclRemoteError(网络错误)FTAR 重试,返回可重试的 work 给上层控制平面,并取消/标记已排队的后续 GPU 操作为 skipped
不可恢复错误ncclSystemError、ncclInvalidArgument(bug 或硬件故障)立即退出崩溃,重试只会浪费时间

FTAR 实测吞吐接近原生 NCCL(见 §七 Table 2)。

5.5 非阻塞追赶(Non-blocking Catch-up)

非阻塞追赶

副本可能因多种原因落后(恢复副本初始化后是随机状态,须从其他副本取状态;或故障后剩余副本状态不一)。

传统追赶难题:落后副本取最新检查点后再执行自检查点以来的任务,但其他副本还在训练新数据,取完检查点后落后副本仍落后。传统上须停顿/减速其他副本,不可接受。

FT-HSDP 的非阻塞追赶协议——利用训练的特殊性质:

每步开始时,各副本通过共识服务上报下一步步号。步号最高(n)的副本视为健康,推进训练 step n;其余副本落后,去健康副本取 step n-1 末对应的检查点。步末,健康副本正常发送梯度,落后副本发送零梯度。梯度交换协议(FTAR)会自动把所有副本带到同一状态。

示例(Figure 6):起初 Replica 3 落后(在 step 96),其余在 step 100。交换信息后 Replica 0-2 知道自己健康、推进 step 100;Replica 3 发现落后,去取 step 99 对应的检查点。步末 Replica 0-2 发梯度、Replica 3 发零梯度。FTAR 把所有副本带到同一状态,随后一起训练 step 101。

两个使其可行的训练系统性质:

  1. 零梯度性质:只要所有副本有相同检查点,一个不做训练工作的副本发送零梯度就能与做了训练工作的副本达到相同状态(传统复制系统中落后副本必须真正执行任务)。
  2. 快速检查点性质:得益于高效检查点恢复(§5.6),取检查点通常比一步还快,健康副本无需等待落后副本。(若取检查点比训练一步慢会导致阻塞,但经验中极罕见。)

5.6 高效检查点与恢复

检查点获取

常态检查点:GPU 每 100 步将自身状态写入远程持久存储。状态含模型状态、优化器状态、dataloader 状态。

  • 模型状态和优化器状态跨副本复制,只需一个副本写入持久存储。
  • dataloader 状态(标记待训练数据位置)不复制但很小,FT-HSDP 每步把所有 GPU 的 dataloader 状态写入单个 loader_state 文件。

追赶时的点对点获取(Figure 7):持久检查点用于全系统重启;但追赶协议中,恢复副本直接从健康副本获取状态(对应最新步)。因状态已复制,恢复副本只需从其中一个获取,且 FT-HSDP 以负载均衡方式组织获取:健康 GPU 将状态写入 CPU 内存,恢复 GPU 通过 HTTP 协议获取(与前反向通信走的 GPU 高速网络无竞争)。恢复副本从 loader_state 文件恢复 dataloader 状态。

此过程假设:一个 GPU 在另一 GPU 获取其状态时不会改变状态——这在优化器步之前成立。通用分布式系统若要边执行新任务边获取状态,通常需 Copy-On-Write (CoW) 等复杂技术。


六、实现细节

实现要点说明
每个 batch 恰好训练一次ML 专家设定的准则:不跳过数据、不重复训练同一 batch。故障副本借助 loader_state 文件从故障处精确恢复。可能造成末尾 straggler,但经验中不严重(除非某副本持续比他人故障更多)。
降低内存压力为容错希望跑大量小副本(单次故障损失小),但每副本须含全部状态 → 内存压力大。缓解:(1) 测试不同副本数取平衡,当前用 10–20 个副本;(2) 将优化器状态从 GPU 内存 offload 到 CPU 内存,按需取回(GPU 取 1GB 优化器状态约需 60ms)。
降低首步效应尽量让恢复副本在向 FT-HSDP 报到前预先执行耗时初始化(重建 checkpointer、初始化 dataloader)。但加载首批数据必须在报到与首个 FTAR 之间完成,故只能减轻而非完全消除首步效应。
CPU 大规模仿真在 100K GPU 上常态跑测试不现实。开发 shadow 模块加载机制,透明地用 CPU mock 模块替换 GPU 模块(计算跳过、只传张量形状;网络用自研 CPU 通信库替代 NCCL)。虽无法查 GPU 相关代码问题,但可测 dataloading、预处理及外部依赖服务(调度、存储、遥测、日志)。实例:早期 get_quorum(决定哪些副本健康的共识协议)在 100K 规模延迟 9 秒,用此工具迭代降到 700ms,真机验证延迟一致。

七、实验结果

7.1 停机时间与训练效率(98K GPU 生产实验)

设置:98K H100 GPU,分为 12 个副本,跨 4 座 DC 楼宇。每副本用 FSDP + 异步张量并行 + 上下文并行 + 定制流水线并行调度,在 8192 GPU 上训练一个数万亿参数的稠密 transformer(文本+图像 token)。

关键结果:

场景吞吐/停机
12 健康副本(稳态)~450 TFlops/GPU/s,与无 FT-HSDP 时相同 → 无故障时零开销
杀掉一个副本(检测+处理+重训故障步)约 3 分钟(含 60 秒超时 + 数秒重建 FTAR + ~20 秒重训一步)
11 副本继续训练~450 TFlops/GPU/s,符合预期
故障副本在 step n 重新加入 → 完成 step n+1约 2 分钟(约 100 秒停机)
完全恢复后 12 副本~450 TF/GPU/s,完全恢复

3 分钟停机的两个意外来源(超出预期,团队已找到修复):

  1. 功耗平滑器(power smoother):故障导致 98K GPU 全停会触发功耗平滑机制,带来 45 秒停顿。
  2. 一个 bug:导致不必要的 FTAR 重配置,额外 45 秒停顿。 修复后停机应约 1.5 分钟,接近预期(60 秒超时 + 数秒重建 FTAR + ~20 秒重训一步)。

100 秒重新加入停机的两个原因:

  1. 首步效应:即使把部分额外任务前移,仍未完全消除(但即便全部 100 秒都归因于此,也短于 Figure 2 中约 200 秒的首步时间,说明优化有效)。
  2. 步中加入的权衡:若副本在 step n 中途或末尾才恢复,FT-HSDP 需在”让其 step n+1 才加入(无停机但 n+1 不利用该副本)“与”让其他副本在 step n 末等待(有停机但 n+1 可用该副本)“之间权衡,设了等待时限故可能产生停机。

效率提升分析(假设每 18 分钟一次故障):

同步恢复:18−1018=44%\text{同步恢复:} \quad \frac{18-10}{18} = 44\%

FT-HSDP(12 副本,全修复仍需 10 分钟,但只完全停机 3 分钟,其余 7 分钟以 11 副本运行):8+7×111218≈80%\text{FT-HSDP(12 副本,全修复仍需 10 分钟,但只完全停机 3 分钟,其余 7 分钟以 11 副本运行):} \quad \frac{8 + 7 \times \frac{11}{12}}{18} \approx 80\%

进一步可用更快的故障检测器替换 60 秒超时来压缩这 3 分钟停机。

7.2 模型精度

核心问题:频繁增删副本(导致每步训练数据量变化)是否影响精度?

设置:全尺度只能短时跑一次(进度符合预期);更全面研究在 256 H100 GPU 上进行,训练多个 3B 激活参数、16 专家的 MoE 模型,训练集为 5000 亿多样文本 token。

实验命名:每个设置表示为 freq_fp_len_con_lrfreq\_fp\_len\_con\_lr:

  • freqfreq:每 n 步故障次数(最高 lenlen=4k 时 n=11k,其他 n=5k)
  • fpfp:浮点精度(fp8 或 fp16)
  • lenlen:故障时长(步数)
  • concon:每次故障并发杀掉的副本数
  • lrlr:学习率干预策略

例:d2x_fp8_for4k_1reps_lr_noned2x\_fp8\_for4k\_1reps\_lr\_none = 每 11k 步 2 次故障、fp8、每次故障持续 4k 步、每次杀 1 副本、无学习率干预。

学习率干预策略(因不同步训练数据量不同,宜纳入考量):

策略缩放公式效果
nonenone基线按步号调整基线
linearlinear×num_healthy_replicasnum_total_replicas\times \dfrac{num\_healthy\_replicas}{num\_total\_replicas}一般(论文未展示)
sqrtsqrt×num_healthy_replicasnum_total_replicas\times \sqrt{\dfrac{num\_healthy\_replicas}{num\_total\_replicas}}更优,保证学习率始终与梯度噪声标准差成正比

sqrtsqrt 通常优于 linearlinear,故论文不展示 linearlinear 结果。

整体训练损失

整体编码集评估

整体进度(Figure 8,200B token 时开始杀副本,x 轴为已见 token 数):无论是训练损失还是编码评估集(另外三个评估集——推理、数学、通用文本——趋势相同),基线与各 FT-HSDP 设置之间无可区分差异,即使最激进的 d2x_fp8_for4k_3reps_Xd2x\_fp8\_for4k\_3reps\_X 场景亦然。(注:若 x 轴改为运行时间,FT-HSDP 因处理恢复会更慢到达相同进度。)

短故障波动

长故障波动

放大观察波动(Figure 9):

  • (a) 短窗口内出现波动,3 个并发故障比 1 个故障波动更大(符合预期)。可见 sqrtsqrt 学习率干预有效平滑曲线(利于调试)。
  • (b) 长故障下结论仍成立,即便最激进的 d2x_fp8_for4k_3reps_lr_Xd2x\_fp8\_for4k\_3reps\_lr\_X 设置,训练精度最终回到基线。

结论:异步恢复最终不会显著降低训练精度,虽然训练中可能有波动。推荐平方根学习率干预来平滑波动。fp16 或不同评估集下结论一致。

7.3 FTAR 性能(微基准)

Table 2:FTAR 带宽(GB/s,100 次运行平均),对比原生 NCCL:

消息大小2 Ranks (FTAR/NCCL)4 Ranks (FTAR/NCCL)8 Ranks (FTAR/NCCL)16 Ranks (FTAR/NCCL)
256MB40.06 / 40.5341.17 / 41.2040.12 / 41.8241.25 / 41.17
512MB39.80 / 42.1242.15 / 42.3243.01 / 43.6143.61 / 43.15
1GB40.93 / 43.8744.67 / 43.9344.67 / 44.8745.51 / 44.32

结论:FTAR 性能接近原生 NCCL(多数场景差距在个位数百分比内,部分场景甚至反超),证明 CPU-GPU 混合设计在获得复杂控制能力的同时几乎不牺牲吞吐。


八、核心创新

创新点说明依据
副本级容错单元以数据并行副本为容错单元,故障时只重建单个副本,其余继续训练恢复范围从整集群缩小到单副本;有效训练时间 44%→80%
CPU-GPU 混合 FTARCPU 驱动控制平面(动态增删连接、拥塞控制、错误分类),GPU 执行数据平面(RDMA 传输)解决 NCCL 无法动态重配置的根本痛点;吞吐接近 NCCL(Table 2)
非阻塞追赶(零梯度)落后副本取检查点后发送零梯度,靠 FTAR 自动同步,无需停顿健康副本利用训练特有的零梯度性质 + 快速检查点性质
首步开销前移将 checkpointer/dataloader 初始化前移到报到之前100 秒 < 200 秒首步时间
平方根学习率干预按 健康副本数/总副本数\sqrt{健康副本数/总副本数} 缩放学习率,使其正比于梯度噪声标准差平滑异步恢复引入的训练波动(Figure 9)
低占用高带宽 FTAR kernel仅用 2 个 SM(<NCCL 的 4 SM),靠 ILP 与寄存器优化在低占用下打满内存总线最小化与并发计算 kernel 的资源竞争
CPU 大规模仿真shadow 模块用 CPU mock 替换 GPU 模块,在无 GPU 时测试 100K 规模控制逻辑提前把 get_quorum 延迟从 9s 降到 700ms

九、相关工作

  • 容错训练:多数系统用同步检查点恢复(MegaScale、RobustLLM/ByteDance、Unicorn、bloom),面临与本文基线相同的挑战。弹性训练(Horovod、PyTorch Distributed、Singularity)天然支持异步恢复,但需从检查点重启或昂贵的模型状态 re-shuffle,且未考虑大规模下重建 NCCL 链路的长延迟。Google Gemini 报告可在局部故障时用更少 TPU 切片继续训练,但未披露细节。Oobleck 和 ReCycle 用流水线并行重建流水线,但未考虑跨 DC 长延迟与有限带宽。Bamboo 用冗余计算容错,但无故障时也有性能代价。
  • 全归约算法:HPC 与 ML 领域对 ring 及其变体有大量研究。GPU 相关实现从早期 CPU 驱动(GPU 缺 I/O 能力)演进到 GPUDirect 再到 NCCL(几乎完全绕过 CPU)。本文关键贡献:证明 CPU 控制平面 + GPU 数据平面的混合设计能同时获得复杂控制能力和接近 NCCL 的性能。
  • 其他:并行训练(FSDP、Megatron、DeepSeek-V3、FlashAttention)、高效检查点(CheckFreq、Check-N-Run)、大规模系统 mock 仿真、学习率调整、故障根因分析等。

十、总结

核心贡献

  1. 详细分解了 O(100K) GPU 规模同步恢复各步骤耗时,论证了其在该规模的不可行性(有效训练时间仅 44%)。
  2. 提出 FTAR 容错全归约协议:无故障时快速集合通信,故障时灵活控制(CPU 控制 + GPU 数据)。
  3. 提出非阻塞恢复协议(零梯度追赶),副本恢复时最小化停机。
  4. 证明 FT-HSDP 的异步恢复不显著影响模型精度。

关键数据

指标同步恢复FT-HSDP改善
停机时间10 分钟3 分钟(修复后预期 1.5 分钟)~3.3×
有效训练时间44%80%~1.8×
无故障吞吐开销—0%—
FTAR 吞吐 vs NCCL—接近持平—

技术影响

  • 为 O(100K)+ GPU 超大规模、跨数据中心 LLM 训练提供了一套经生产验证的容错范式,将有效训练时间从 44% 提升到 80%。
  • CPU-GPU 混合通信设计的思路可能影响未来容错集合通信库的设计方向——GPU 未必要独占通信控制。

局限性与未来方向

  • 精度结论主要从 256 GPU 小规模外推:全尺度只跑了一次短时实验。作者承认这是资源约束下的标准无奈做法(小规模调参后直接套用到全尺度)。
  • 首步效应未完全消除:重新加入仍有约 100 秒停机。
  • 内存压力:小副本利于容错但每副本须含全部状态,当前折中为 10–20 副本 + 优化器状态 CPU offload。
  • 超时检测偏保守:60 秒超时是停机主因之一,未来计划与快速根因分析集成以加速检测。
  • straggler 问题:“每 batch 恰好训练一次”策略可能造成末尾 straggler;未来可探索故障时 batch 重排等更复杂策略。

十一、参考资源


附录:论文图表索引

图表文件说明
Figure 1anetwork-building.pngDC 楼内网络架构(3 层 Clos)
Figure 1bnetwork-mesh.png跨 DC 全连接网格
Figure 2recovery-time-breakdown.png恢复时间随规模增长的分解
Figure 3architecture.pngFT-HSDP 整体架构
Figure 4rank0-leader-design.pngRank 0(Leader)设计
Figure 5ftar-protocol.pngFTAR 协议细节
Figure 6non-blocking-catchup.png非阻塞追赶协议
Figure 7fetch-checkpoint.png检查点获取流程
Figure 8aaccuracy-loss-overall.png训练损失整体曲线
Figure 8baccuracy-eval-overall.png外部编码集评估
Figure 9aaccuracy-short-failure.png短故障训练进度波动
Figure 9baccuracy-long-failure.png长故障训练进度波动