知识库 / 基础设施
GitHub
← 基础设施

LAYER 01 / INFRASTRUCTURE

分布式存储、数据读取与 Checkpoint

所属: 基础设施层。本页边界: 训练数据和模型状态怎样以足够吞吐被读取、保存和恢复。

#核心公式

TI/O≥Sread/writeBend−to−end,τ∗≈2Cμ T_{\mathrm{I/O}}\geq\frac{S_{\mathrm{read/write}}}{B_{\mathrm{end-to-end}}}, \qquad \tau^*\approx\sqrt{2C\mu}

第一式给出数据读写时间下界;第二式是在简化故障模型下,使检查点开销与故障丢失计算达到平衡的近似最优保存间隔。

#符号说明

符号 含义 单位/条件
TI/OT_{\mathrm{I/O}} 一次读写操作的端到端时间 s
Sread/writeS_{\mathrm{read/write}} 本次读取或写入的数据量 Byte
Bend−to−endB_{\mathrm{end-to-end}} 存储到消费者整条路径的有效带宽 Byte/s
τ∗\tau^* 近似最优 Checkpoint 间隔 s
CC 完成一次 Checkpoint 的时间 s
μ\mu 平均无故障运行时间 s
R, P, CgR,\ P,\ C_g 后文中的读取、处理、设备消费速率 sample/s 或 Byte/s
BeffectiveB_{\mathrm{effective}} 数据链路各阶段共同限制下的端到端吞吐 Byte/s
Bstore, Bnetwork, Bdecode, Bhost→deviceB_{\mathrm{store}},\ B_{\mathrm{network}},\ B_{\mathrm{decode}},\ B_{\mathrm{host\to device}} 存储、网络、解码与主机到设备阶段吞吐 Byte/s
Bpipeline, BjB_{\mathrm{pipeline}},\ B_j 整条流水线吞吐、第 jj 个处理阶段吞吐 Byte/s 或 sample/s,口径需一致
H(τ)H(\tau) 单位有效运行时间的近似额外开销比例 无量纲

τ∗\tau^* 公式假设故障近似独立、保存时间稳定,并忽略恢复和更高阶项;实际系统需用实测故障与保存成本校准。

#技术要点

  • 对象存储、分布式文件系统和本地 NVMe 的延迟/吞吐不同。
  • 并行文件系统(Lustre、GPFS)为顺序大 I/O 优化,元数据操作和随机小文件是其主要弱点。
  • 数据分片与顺序读取降低小文件和元数据开销;小文件随机读取是训练数据加载的核心故障模式。
  • checkpoint 需包含权重、优化器、随机状态和分片布局。
  • 恢复时的重分片是张量布局问题,不是简单复制文件。
  • PyTorch Distributed Checkpoint(DCP)通过逻辑张量抽象支持跨并行配置的重分片。
  • 大规模训练的故障率远高于单机;LLaMA3 在 54 天预训练中经历了 466 次中断。

#原理与演进

#数据读取:吞吐链路

Beffective=min⁡(Bstore,Bnetwork,Bdecode,Bhost→device) B_{\mathrm{effective}}=\min(B_{\mathrm{store}},B_{\mathrm{network}},B_{\mathrm{decode}},B_{\mathrm{host\to device}})
  • 存储带宽足够仍可能被小文件元数据访问、解压、token 解码或主机到设备传输卡住。
  • 大分片顺序读、预取和并行解析提高持续吞吐;随机打乱应与顺序读取做权衡。
  • 训练等待数据会让加速器空转,因此应看有效 token/s,而非只看磁盘 GB/s。实测表明,端到端训练时间中 10%–70% 可能消耗在 I/O 停顿上。

#Checkpoint 的状态闭包

  • 仅推理:权重、结构、tokenizer 与必要配置。
  • 继续训练:还需优化器状态、学习率调度、随机数状态、数据采样位置与并行分片信息。
  • 分片 Checkpoint 减少单节点写入瓶颈,但换并行度恢复时需重新分片;只复制文件不能保证状态一致。
  • 保存间隔的取舍:更频繁保存降低故障后重算量,却消耗更多 I/O 与同步时间。

参见: 数据版本与血缘定义“读的是哪一版数据”。

#1. 数据从存储走到训练芯片

持久化文件 → 元数据定位 → 顺序/随机读取 → 解压与解析 → 主机内存预取 → 设备传输 → token batch。

  • 任一环节供给速度低于训练消费速度,就会让计算设备等待。
  • 小文件过多时,瓶颈可能是打开文件与元数据查询,而非字节带宽;把样本组织为可顺序读取的分片通常更有利。
  • 压缩减少磁盘与网络字节数,却增加 CPU 解压;最优压缩率取决于哪一段链路最慢。
  • 预取通过流水化隐藏读取延迟,但需要主机内存缓存;缓存容量不足会导致频繁重新读取。

#1.1 存储形态对比:并行文件系统、对象存储与本地 NVMe

大规模训练集群的存储选择直接影响数据加载和 Checkpoint 性能。三种主要形态各有明确的适用边界:

并行文件系统(Lustre、GPFS/Spectrum Scale)。 这类系统为 HPC 的顺序批量 I/O 设计,通过元数据服务器(MDS)和对象存储服务器(OSS)分离元数据与数据路径,支持数百 GB/s 的聚合带宽。Lustre 是 DGX SuperPOD 参考架构的标准选择。但其根本弱点是元数据操作和小文件随机读取:Lustre 的 MDS 在大量小文件操作时容易成为瓶颈,基于 HDD 的 Lustre 随机读取性能受限于寻道时间。GPFS 在元数据管理上比 Lustre 更成熟(客户端节点参与元数据管理),但架构复杂度和运维成本更高。

对象存储(S3、OSS、MinIO)。 对象存储为容量和持久性优化,而非延迟。每次按样本读取对象都产生 per-object 延迟;直接在训练中逐样本访问 S3 会导致严重的延迟瓶颈。实践中需要在对象存储前加缓存层(如 FSx for Lustre 缓存 S3,或本地 NVMe 缓存),将冷数据访问延迟降低一个数量级。

本地 NVMe。 PCIe Gen5 NVMe SSD 的顺序读取吞吐已超过 14 GB/s,是训练层存储的标准配置。节点本地 NVMe 的延迟最低,但容量有限,且数据不跨节点共享。典型策略是:将训练数据预分片到节点本地 NVMe,用分布式采样器确保每个 rank 读取自己的分片。

分层策略。 数据放置应尽可能靠近计算:节点本地 NVMe → 机架本地 NVMe-oF → 并行文件系统。如果数据集能放入主机 RAM,启动时预加载并完全跳过磁盘是延迟最低的方案。

#1.2 数据路径的软件开销:io_uring、SPDK 与 GPUDirect Storage

存储硬件的标称带宽与实际可达带宽之间的差距,往往来自软件 I/O 栈的开销。在 LLM 训练场景中,三种数据路径的适用性有明确差异:

io_uring。 Linux 内核提供的异步 I/O 接口,通过共享环形缓冲区(submission queue 和 completion queue)在用户态和内核态之间传递 I/O 请求,减少了传统 libaio 的系统调用和上下文切换开销。在推理场景的小随机 I/O 中,io_uring 实现了最低延迟和具有竞争力的 IOPS。

SPDK(Storage Performance Development Kit)。 用户态 NVMe 驱动,通过轮询模式彻底绕过内核 I/O 栈。SPDK 的局限是缺乏 POSIX 文件系统支持,只能操作裸块设备。适用于需要极致吞吐且能接受自己管理数据布局的场景。

GPUDirect Storage(GDS)。 允许 NVMe SSD 通过 PCIe 直接与 GPU 显存进行 DMA 传输,绕过主机内存和 CPU。在预训练和微调场景中,负载以粗粒度顺序读写为主,GDS 在降低加载时间和主机 CPU 使用率方面表现最佳。在 CPU 介导的数据路径中,每核心 GB/s 是区分不同方案的关键指标。设计原则是:推理用 io_uring 优化小随机 I/O 延迟,预训练/微调用 GDS 提高每核心吞吐。

#1.3 小文件随机读取:核心故障模式

存储对大顺序读取的吞吐远高于小随机读取。数百万个独立样本文件(每张图片一个 JPEG、每篇文档一个 JSON)使每个 epoch 变成 seek storm,GPU 越快,这个瓶颈越早显现。

修复方式在离线阶段完成,而非运行时:将大量样本分片为少数大文件——Arrow、Parquet、TFRecord、WebDataset tar,或 NeMo 的 memory-mappable .bin/.idx 对。一次 chunked read 产生多个样本,将 I/O 模式从随机变为顺序。对象存储同理:训练开始前将小 S3 对象合并为大对象。

WebDataset 的机制。 WebDataset 使用 POSIX tar 格式存储分片,每个分片包含数千个样本。它设计为顺序只读:数据加载器打开一个分片,顺序读取样本,然后移动到下一个分片。这种模式充分利用了存储的顺序带宽,同时避免了文件打开/关闭的系统调用风暴。在 PyTorch DataLoader 中,每个 worker 分配唯一分片,num_workers 和 prefetch_factor 控制并发预取深度。

诊断表。

症状 可能原因 处理方向
IOPS 高但 MB/s 低 数百万小文件 分片为大的顺序文件
读取大小约 4 KB 未调优的 buffer/prefetch chunk 将读取块提升到约 1 MB
随机访问确实不可避免 每次系统调用开销占主导 并行 pread() 线程或 io_uring
网络被冗余读取饱和 每个节点读取整个数据集 预分片到节点本地 NVMe;每 rank 用 DistributedSampler
对象存储逐样本读取 per-object 延迟 暂存到本地 NVMe,或加缓存层

#1.4 NUMA 与拓扑感知的数据加载

数据加载线程与目标 GPU 的 NUMA 亲和性直接影响主机到设备传输带宽。在 8 卡服务器中,CPU 0 管理 GPU 0–3 和 NIC 0–1,CPU 1 管理 GPU 4–7 和 NIC 2–3。如果 GPU 0 的数据处理线程运行在 CPU 1 上,数据需要跨 CPU 互联(Intel UPI 或 AMD Infinity Fabric)传输,延迟和带宽都变差。

正确做法是让 pin_memory 线程与目标 GPU 位于同一 NUMA 节点,并使 DataLoader 的 worker 进程也绑定到本地 CPU。HuggingFace 和 WebDataset 的 DataLoader 实现都支持通过环境变量或进程绑定来控制亲和性。在评测数据加载性能时,应测量每个 rank 的有效样本消费速率,而非仅报告存储端的聚合带宽。

#2. 吞吐的上下界与排队

设各环节长期吞吐为 B1,…,BkB_1,\ldots,B_k:

Bpipeline≤min⁡jBj. B_{\mathrm{pipeline}}\leq\min_j B_j.
  • 这是上界;同步边界、小批次和尾部等待可能让实际吞吐更低。
  • 加快已经不是瓶颈的一环,对端到端 token/s 作用有限。
  • 数据打乱不能只在文件名层面做:如果同类样本在大分片内部集中,训练批次仍可能高度相关。
  • 全局逐样本随机读取会破坏顺序 I/O;常用分片级打乱加局部缓冲打乱折中。

#3. Checkpoint 具体要保存哪些状态

状态 只做推理 精确继续训练时
模型权重与结构配置 必需 必需
Tokenizer 与消息模板 必需 必需
优化器动量与主权重 不需要 通常必需
学习率调度、步数与 loss scaler 不需要 使用时必需
随机数状态与采样位置 不需要 对可复现继续训练重要
分片/并行布局元数据 视制品格式 恢复与重分片时重要
  • 只保存权重可以继续做微调,但不能称为“从同一点无缝恢复原训练动态”:优化器历史与数据顺序已改变。
  • 分布式训练把状态分散在设备上,Checkpoint 也常被分片写出;更换设备数时需要按逻辑张量重新组装和分片,不是按原文件名机械复制。

#3.1 PyTorch Distributed Checkpoint 的机制

PyTorch Distributed Checkpoint(DCP)是当前分布式 Checkpoint 的主要实现,其核心抽象是逻辑张量(logical tensor) 与物理分片(physical shard) 的分离。

保存时,每个 rank 提供自己的本地 state_dict,DCP 的 SavePlanner 根据分片元数据将本地张量映射到逻辑张量的对应切片,并写入存储后端。加载时,LoadPlanner 根据目标 state_dict 请求的布局读取并组装切片。只要保存端和加载端对逻辑张量的名称、形状与分片语义保持兼容,DCP 就能在加载时执行重分片(resharding);例如把保存时的 TP=4 布局转换为目标端能够正确描述的 TP=8 布局。它并不保证任意并行配置之间都能无条件转换。

PyTorch 的统一 get_state_dict 接口可为受支持的 DDP、FSDP 和张量并行组合生成一致的逻辑状态;流水线并行、优化器状态以及框架自定义布局仍可能需要额外元数据、规划器或显式转换。Megatron-LM 的 sharded_state_dict 也只会重建框架明确描述过的逻辑分片,不能把 TP、PP、DP、EP、CP、FSDP 的任意变化都视为自动兼容。Fairseq2 的 reshard_tensor 展示了底层思路:先依据已知分片语义恢复逻辑张量,再按目标布局重新切片。

#3.2 异步保存的内存开销与一致性

DCP 的异步保存(async_save)将分片的 state_dict 先复制到 CPU 缓冲区,然后在后台线程中写入持久化存储,从而不阻塞训练。这一机制的必要性在于:在 Checkpoint 写入期间,模型和优化器权重不能继续变化,否则保存的状态不是一致的。

代价是 CPU 内存增加约 checkpoint_size_per_rank × number_of_ranks。对于千亿参数模型,所有 GPU 的 Checkpoint 内存优先保存在主机内存,建议预留 200 GB 以上。

异步保存的一致性前提是:在开始 CPU 拷贝的时刻,所有 rank 的训练状态处于同一逻辑步。如果不同 rank 在不同步保存权重和优化器,得到的组合状态可能从未在真实训练中存在。因此,分布式 Checkpoint 的同步点应先确定同一逻辑步的状态,再各自异步写入。

#4. 一致性与异步保存

  • 若不同设备在不同训练步保存权重和优化器,得到的组合状态可能从未在真实训练中存在。
  • 同步保存先确定同一逻辑步的状态,再写持久化介质;等待 I/O 会中断计算。
  • 异步保存可缩短计算停顿,但需要保留一致的快照、管理写入缓冲,并处理未完成写入时的故障。
  • Checkpoint 成功不只意味着文件存在,还需要能识别完整分片与版本;这一点是状态一致性的必要条件。

#5. 保存频率的数学取舍

设每次保存耗时 CC,两次保存之间的计算时间为 τ\tau,平均故障间隔为 μ\mu。在简化的独立故障模型下,长期额外时间占比近似:

H(τ)≈Cτ+τ2μ. H(\tau)\approx\frac{C}{\tau}+\frac{\tau}{2\mu}.

对上式关于 τ\tau 求极小值,即得到页首公式,不在此重复列式。

τ∗≈2Cμ \tau^*\approx\sqrt{2C\mu}
  • 第一项:保存过于频繁,I/O 开销大。
  • 第二项:保存太少,故障后平均需重算约半个间隔。
  • 该式忽略恢复时间、异步重叠、多故障与非平稳风险;它解释为什么有最优间隔,不是通用配置值。

#5.1 大规模训练的故障现实与连续 Checkpoint

LLaMA3 预训练在 54 天内经历 466 次中断,平均 MTBF 约 2.8 小时。OPT-175B 训练中平均每 2.5 小时出现一次硬件故障。在这种故障率下,固定间隔的 Checkpoint 面临两难:间隔太长则故障后重算代价高,间隔太短则 I/O 开销累积。

Google 的 Orbax 和 MaxText 引入了连续 Checkpoint(Continuous Checkpointing) :不再按固定步数或时间触发,而是在上一次保存成功完成后立即异步发起下一次保存。这一策略最大化利用了可用 I/O 带宽——只要存储能跟上,Checkpoint 就以存储的极限速率持续进行,同时最小化故障后丢失的计算量。minimum_interval_secs 参数用于设置冷却期,避免在极短时间内产生过多 Checkpoint。

动态间隔调整进一步优化这一过程:使用抢占预测和 Checkpoint/恢复代价的数学模型,在可抢占云 VM 上动态调整间隔,比固定间隔训练吞吐提升最高 60%。

#5.2 分层与容错恢复策略

TierCheck 提出三层 Checkpoint 设计,将存储放置与故障异质性对齐:Tier-1 在 GPU 内存中保存轻量差分 Checkpoint,Tier-2 在本地 NVMe 和 peer 内存中保存,Tier-3 在远程持久化存储中保存。软件和节点故障优先从 Tier-1/Tier-2 快速恢复,避免从远程存储加载全部状态。

CPSR 放弃了“重算”路径:用轻量预测器根据常规 Checkpoint 预测故障前的训练状态,预测误差作为微小扰动由训练过程自行纠正。在相同测试环境下,恢复代价比现有方案降低 41.66%,GPU 内存占用增加不足 200 MB。

Resilio 针对千亿参数模型,通过多层次优化 Checkpoint 读写和即时保存机制,将故障恢复时间缩短至 10 分钟以内。

#6. 存储形态的演进

单机文件 → 分布式共享存储 → 大规模顺序分片与缓存 → 分片、异步、可重分布的 Checkpoint。

解决的旧问题 新技术 新问题
单机容量与共享受限 分布式存储 网络与元数据争用
读小文件开销大 顺序分片 打乱粒度下降
保存阻塞训练 异步快照 一致性与缓冲内存
并行度变化难恢复 逻辑张量重分片 元数据与读写规划复杂
固定间隔在故障率下失效 连续 Checkpoint + 动态间隔 需要准确的故障/代价模型
远程存储加载慢 分层 Checkpoint 层间一致性和管理复杂度

数据逻辑版本在血缘模块;这里只讨论物理读取和状态持久化。

#7. 存储栈调优的工程边界

块层。 现代 Linux 使用 blk-mq 多队列调度器。NVMe 的正确设置是 none(低延迟默认)或 mq-deadline;旧的 CFQ 调度器已废弃。流式大文件时,通过 blockdev --setra 将 read_ahead_kb(默认约 128 KB)提升到数 MB。

文件系统。 XFS 是 Linux NVMe 的常见选择,挂载时加 noatime 消除每次读取的访问时间写入。NFS 仅在少量节点时可行;单 NFS 服务器在多节点并发读取时很快成为吞吐瓶颈,并行文件系统或云缓存是替代方案。

RAID 与 PCIe 通道。 确认 SSD 位于足够的 PCIe 通道上;单块盘无法饱和 GPU 消费速率时,用 RAID 0 条带化多块 SSD。PCIe Gen5 NVMe 的顺序读取速度超过 14 GB/s,但只有在数据路径的软件开销被最小化时才能达到。

测量口径。 数据加载性能必须用有效样本消费速率(samples/s 或 tokens/s)衡量,而非存储端的聚合带宽。同时记录每个 rank 的 I/O 等待时间、CPU 利用率(解压和 tokenization 是否成为瓶颈)和主机到设备传输的 P99 延迟。

#原始资料

本页由仓库中的 Markdown 生成。具体技术结论请结合正文引用与实验条件理解。

输入关键词,探索整个知识库

↑ ↓ 选择 ↵ 打开36 篇笔记,一次搜索