1. 大模型训练里最容易被低估的环节:数据管道
做过大模型预训练的人都有一个共识:模型结构决定上限,数据质量决定下限,而数据管道的效率直接决定你多久能摸到这个下限。我见过太多团队在模型并行、算子优化上砸了几周时间,结果训练吞吐上不去,最后定位下来问题出在 DataLoader 上——GPU 利用率长期在 60% 上下晃荡,算力全耗在等数据了。
MindSpore Transformers 这套框架里的 Blended Megatron DataLoader,就是专门解决这个问题的。它做的事情说起来不复杂:把多个不同来源、不同格式的数据集按指定权重混合成一个统一的训练数据流,同时保证在分布式训练场景下每个卡拿到的数据既不重复也不遗漏,还要让数据预取和计算重叠起来,把 I/O 等待藏到计算背后。
这套机制适合谁看?如果你正在用 MindSpore 做 LLM 预训练或者微调,需要自己准备和处理训练数据,或者你发现训练时 step time 波动大、GPU 利用率上不去,那这篇文章里的内容应该能帮到你。即便你用的是别的框架,Megatron 系的数据加载思路也是通用的,理解了这个设计,换个框架照样能迁移。
我下面会从整体设计思路开始拆,然后逐层深入到配置细节、实操步骤、参数计算,最后把我踩过的坑和排查经验整理出来。内容偏实操,代码和配置都会给到可以直接参考的版本。
2. 整体设计思路:为什么要做 Blended DataLoader
2.1 从单数据集到多数据集混合的需求演变
早期做小模型训练的时候,数据加载很简单:一个 Dataset 对象,一个 DataLoader 包一层,设置好 batch_size 和 shuffle 就完事了。但到了 LLM 预训练阶段,情况完全变了。
首先,预训练数据通常不是单一来源。你可能需要混合网页爬取数据、书籍语料、代码数据、学术论文等多个来源,而且每个来源的权重还不一样。比如网页数据占 60%,代码数据占 20%,书籍和论文各占 10%。这种加权混合的需求,用简单的 ConcatDataset 是满足不了的,因为 ConcatDataset 只是把数据集首尾相接,不提供权重控制。
其次,数据格式不统一。有的数据是 jsonl 格式,每行一个 json 对象;有的是 bin 格式的二进制 token 序列;有的还是原始的 txt 文本需要在线 tokenize。如果每个格式都单独写一个 DataLoader,然后手动控制采样比例,代码会变得非常难以维护。
第三,分布式训练带来的复杂性。当你用几十张卡甚至上百张卡做数据并行的时候,必须保证每张卡在每个 step 拿到的数据是全局唯一的,否则相当于变相减小了有效 batch size,浪费算力。同时还要保证数据加载不会成为瓶颈,需要预取机制。
Blended Megatron DataLoader 的设计目标就是一次性解决这三个问题:权重混合、格式统一、分布式高效加载。
2.2 Megatron 风格数据加载的核心设计哲学
Megatron 系的数据加载有一个非常核心的设计理念:索引化(Indexed Dataset)。这个思路值得展开说一下。
传统的数据加载是流式的:从文件头读到文件尾,一条一条往外吐。这种方式在单机单卡场景下没问题,但在分布式场景下就很麻烦——你没法让第 3 号卡直接跳到第 15000 条数据开始读,只能从头遍历。
索引化的做法是:预先为每个数据文件建立一个索引文件,记录每条数据的偏移量和长度。这样加载的时候,给定一个全局索引号,就能直接定位到对应的数据位置,实现随机访问。这带来的好处是:
- 分布式切分变得简单:全局索引按 rank 和 world_size 取模分配即可
- 断点续训容易实现:记录当前消费到的索引位置,恢复时直接从该位置继续
- 数据打乱灵活:只需要打乱索引数组,不需要移动实际数据
MindSpore Transformers 的 Blended Megatron DataLoader 正是基于这个思路构建的。它把多个数据集的索引合并成一个全局索引表,然后按照权重进行采样分配,最后通过分布式采样器把索引分配到各个卡上。
2.3 权重混合策略背后的考量
混合权重的设计看起来简单,但实际使用中有几个容易踩坑的地方。
第一个坑是权重归一化。你配置的权重不一定是加起来等于 1 的,比如你写 [3, 1, 1],实际含义是 3/5、1/5、1/5。框架内部会做归一化处理,但你需要清楚最终每个数据集被采样的概率是多少。
第二个坑是数据量不匹配导致的重复采样。假设数据集 A 有 100 万条,数据集 B 只有 1 万条,权重各占 50%。那么训练过程中 B 会被反复采样很多遍,而 A 可能一轮都没走完。这本身不是 bug,但如果你没意识到这一点,可能会对训练效果产生困惑。实际使用中,通常建议权重配置和数据量大致匹配,或者对小数据集做适当的上采样。
第三个坑是采样粒度和 epoch 边界。Blended 采样是在样本级别进行的,不是数据集级别。也就是说,每个 batch 里可能同时包含来自不同数据集的样本。这跟"先训完 A 再训 B"的课程学习策略是完全不同的。如果你需要课程学习,得用另外的机制。
3. 核心细节解析:从配置到执行的完整链路
3.1 数据格式与索引构建
MindSpore Transformers 支持的数据格式主要有两种:MindRecord和Megatron 二进制格式。两者各有适用场景。
MindRecord 是 MindSpore 原生的数据格式,优点是 schema 灵活、支持多种数据类型、自带索引。缺点是写入和读取的开销相对较大,对于超大规模 token 序列数据,存储效率不如纯二进制格式。
Megatron 二进制格式则是专门为 LLM 训练优化的。它由两个文件组成:.bin文件存储实际的 token 序列(通常是 uint16 或 uint32),.idx文件存储索引信息(每条数据的起始偏移和长度)。这种格式的读取效率极高,几乎就是内存映射加指针跳转,非常适合大规模预训练。
构建索引的过程通常是在数据预处理阶段完成的。以 Megatron 格式为例,核心逻辑是:
import numpy as np def build_index(token_ids_list, output_prefix): """将 token 序列列表写入 bin 文件并构建 idx 索引""" # 写入 bin 文件 all_tokens = np.concatenate(token_ids_list).astype(np.uint16) all_tokens.tofile(f"{output_prefix}.bin") # 构建索引:记录每条数据的起始位置和长度 offsets = np.zeros(len(token_ids_list) + 1, dtype=np.int64) for i, tokens in enumerate(token_ids_list): offsets[i + 1] = offsets[i] + len(tokens) # 索引文件格式:前两个值存维度信息,后续存偏移 header = np.array([len(token_ids_list), 1], dtype=np.int32) index_data = np.concatenate([header, offsets[:-1].astype(np.int32), offsets[1:].astype(np.int32)]) index_data.tofile(f"{output_prefix}.idx")这段代码是简化版,实际框架里的实现会更复杂一些,会处理 dtype 转换、多文件分片、压缩等细节。但核心思路就是这样:bin 文件存数据,idx 文件存偏移,读取时通过偏移直接定位。
注意:构建索引时一定要确保 token 序列的 dtype 和训练时配置的 dtype 一致。我见过有人预处理时用了 uint32,训练配置里写的是 uint16,结果读出来的 token 全是乱的,排查了大半天才发现是类型不匹配。
3.2 分布式采样器的工作机制
分布式采样器是 Blended DataLoader 里最核心也最容易出问题的组件。它的职责是:在全局索引空间上做采样,然后把采样结果分配到各个 rank 上,保证不重不漏。
具体的工作流程是这样的:
第一步,根据各数据集的权重和大小,计算出全局采样序列。假设数据集 A 有 N_A 条,权重 w_A,数据集 B 有 N_B 条,权重 w_B。总采样步数为 T 时,从 A 采样的数量约为 T * w_A / (w_A + w_B),从 B 采样的数量约为 T * w_B / (w_A + w_B)。
第二步,对每个数据集内部进行 shuffle,然后按照计算出的采样数量取出对应的索引。
第三步,将来自不同数据集的索引混合在一起,再次 shuffle,形成全局采样序列。
第四步,将全局采样序列按 rank 切分。如果 world_size 为 W,global_batch_size 为 B,那么每个 rank 每个 step 拿到的数据量为 B/W。切分方式是:第 i 个全局 batch 的第 j 个样本分配给 rank (i * B + j) % W。
这里有一个关键参数:global_batch_size 必须能被 world_size 整除。否则切分时会出现不均衡,某些 rank 会多拿一个样本,导致训练时 step 不一致。框架通常会做检查并报错,但你在配置的时候就要注意这一点。
3.3 数据预取与流水线重叠
数据预取是隐藏 I/O 延迟的关键手段。基本思路是:在 GPU 计算当前 batch 的同时,CPU 侧已经在准备下一个 batch 的数据了。这样当 GPU 算完当前 batch 需要新数据时,数据已经就绪,不需要等待。
MindSpore 里通过dataset.prefetch()或者 DataLoader 的prefetch_size参数来控制预取深度。预取深度设多少合适?这取决于你的 I/O 速度和计算速度的比值。
如果 I/O 很慢(比如从网络存储读取),计算很快,那预取深度要大一些,比如 5 到 10,才能把 I/O 延迟藏住。如果 I/O 很快(比如数据全在内存里),计算是瓶颈,那预取深度设 2 到 3 就够了,设太大反而浪费内存。
我的一般建议是:先用默认值跑一下,观察 GPU 利用率。如果利用率稳定在 90% 以上,说明预取够了。如果利用率波动大,经常掉到 70% 以下,那就加大预取深度。但也要注意内存占用,预取深度乘以单 batch 数据量就是额外的内存开销。
4. 实操过程:从零搭建一个 Blended DataLoader
4.1 数据准备与预处理脚本
假设我们有三份原始数据:web_data.jsonl、code_data.jsonl、book_data.jsonl。每行是一个 json 对象,包含 "text" 字段。我们需要把它们转换成 Megatron 二进制格式。
第一步是 tokenize。这里用 MindSpore Transformers 自带的 tokenizer:
from mindformers import AutoTokenizer tokenizer = AutoTokenizer.from_pretrained("llama2_7b") def tokenize_file(input_path, output_path, max_seq_len=2048): """将 jsonl 文件 tokenize 并写入二进制文件""" all_tokens = [] with open(input_path, 'r', encoding='utf-8') as f: for line in f: data = json.loads(line) text = data.get("text", "") if not text: continue tokens = tokenizer.encode(text) # 添加 EOS token tokens = tokens + [tokenizer.eos_token_id] # 按 max_seq_len 切分 for i in range(0, len(tokens), max_seq_len): chunk = tokens[i:i + max_seq_len] if len(chunk) >= 64: # 过滤过短的序列 all_tokens.append(chunk) return all_tokens这里有几个实操细节值得注意。EOS token 的添加是必须的,否则模型学不会在合适的位置停止生成。序列切分时,最后一段如果太短(比如小于 64 个 token),建议直接丢弃,因为过短的序列对训练贡献很小,反而会增加 padding 开销。
第二步是构建索引并保存:
def save_megatron_format(token_chunks, output_prefix): """保存为 Megatron 二进制格式""" # 展平所有 token flat_tokens = np.concatenate(token_chunks).astype(np.uint16) flat_tokens.tofile(f"{output_prefix}.bin") # 构建索引 lengths = np.array([len(c) for c in token_chunks], dtype=np.int32) offsets = np.zeros(len(lengths) + 1, dtype=np.int64) offsets[1:] = np.cumsum(lengths) # 写入索引文件 with open(f"{output_prefix}.idx", 'wb') as f: # 头部:样本数、维度 f.write(np.array([len(lengths), 1], dtype=np.int32).tobytes()) # 偏移数组 f.write(offsets[:-1].astype(np.int32).tobytes()) # 长度数组 f.write(lengths.tobytes())4.2 配置文件编写与参数说明
数据准备好之后,需要写训练配置文件。MindSpore Transformers 的配置文件通常是 YAML 格式。以下是一个 Blended DataLoader 的配置示例:
train_dataset: type: BlendedMegatronDataset data_path: - /path/to/web_data - /path/to/code_data - /path/to/book_data weights: [0.6, 0.2, 0.2] seq_length: 2048 global_batch_size: 64 shuffle: True seed: 42 num_samples: 1000000 prefetch_size: 4 num_parallel_workers: 8逐项说明关键参数:
weights控制各数据集的采样比例。这里 web 数据占 60%,code 和 book 各占 20%。注意权重列表的长度必须和 data_path 的长度一致。
seq_length是序列长度,必须和模型配置里的 seq_length 一致。如果数据预处理时切分的长度和这里不一致,会出现读取错误。
global_batch_size是全局 batch size,必须能被 world_size 整除。比如 64 的 global_batch_size 在 8 卡训练时,每卡 batch size 为 8。
num_samples是总采样步数。这个值决定了训练一个 epoch 会消费多少条数据。通常设置为数据集总大小的若干倍,具体取决于你想训练多久。
prefetch_size是预取深度,前面已经讨论过。
num_parallel_workers是数据加载的并行线程数。一般设置为 CPU 核数的 1/4 到 1/2。设太大反而会因为线程切换开销导致性能下降。
4.3 启动训练与验证数据流正确性
配置写好后,启动训练的命令通常是:
bash scripts/run_distribute_train.sh 8 configs/llama2_7b_pretrain.yaml启动之后,怎么验证数据流是正确的?我一般会做三个检查。
第一个检查:打印每个 rank 第一个 step 拿到的数据。确认不同 rank 拿到的数据确实不同,而且 token 值在合理范围内(比如在 vocab_size 之内)。
第二个检查:观察 loss 曲线。如果数据流有问题,比如不同 rank 拿到了重复数据,loss 会下降得异常快然后很快过拟合。如果数据混合比例不对,loss 的下降模式也会和预期不符。
第三个检查:用小规模数据跑一个完整 epoch,统计各数据集实际被采样的次数,和配置的权重做对比。偏差应该在合理范围内(比如 5% 以内)。
5. 常见问题与排查技巧实录
5.1 数据加载相关的典型报错与解决
在实际使用中,我遇到过不少数据加载相关的问题。下面整理一个速查表:
| 报错信息 | 可能原因 | 解决方法 |
|---|---|---|
| IndexError: index out of range | 索引文件损坏或与 bin 文件不匹配 | 重新构建索引,确保 bin 和 idx 文件对应 |
| ValueError: seq_length mismatch | 预处理时的序列长度和配置不一致 | 检查预处理脚本和配置文件的 seq_length |
| RuntimeError: batch size not divisible | global_batch_size 不能被 world_size 整除 | 调整 global_batch_size 或 world_size |
| 训练 loss 为 NaN | token 值超出 vocab_size 范围 | 检查 tokenizer 的 vocab_size 和实际 token 值 |
| GPU 利用率低 | 预取深度不够或 I/O 瓶颈 | 增大 prefetch_size,检查存储带宽 |
5.2 性能调优:让 GPU 不再等数据
数据加载性能调优的核心目标是让 GPU 利用率稳定在高位。我的一般调优步骤是这样的:
先看基线。用默认配置跑 100 个 step,记录平均 step time 和 GPU 利用率。如果 GPU 利用率已经在 95% 以上,那基本没什么可调的,瓶颈在计算侧。
如果 GPU 利用率低于 90%,先加大 prefetch_size。从 2 加到 4,再到 8,观察利用率变化。如果加到 8 之后利用率没有明显提升,说明瓶颈不在预取深度。
然后检查 num_parallel_workers。这个参数控制数据加载的并行度。如果 CPU 核数足够,可以适当加大。但要注意,MindSpore 的数据加载线程和计算线程是共享 CPU 资源的,设太大反而会拖慢计算。
再检查存储 I/O。用 iostat 或者类似的工具看看读取带宽是否打满了。如果存储带宽是瓶颈,考虑把数据拷贝到本地 SSD,或者用内存文件系统。
最后检查数据格式。如果用的是 MindRecord 格式,读取开销会比 Megatron 二进制格式大不少。在超大规模训练场景下,建议统一用二进制格式。
5.3 分布式场景下的数据一致性保证
分布式训练里最怕的就是数据不一致。比如某个 rank 挂了重启后,数据流的位置和其他 rank 对不上,导致训练崩溃或者效果异常。
保证一致性的关键是:所有 rank 使用相同的随机种子和相同的采样逻辑。具体来说:
- seed 参数必须在所有 rank 上一致
- shuffle 的逻辑必须是确定性的,不能依赖运行时的随机状态
- 断点续训时,需要保存和恢复采样器的状态,包括当前 epoch、当前 step、随机数生成器的状态
MindSpore Transformers 的 BlendedMegatronDataset 内部已经处理了大部分一致性逻辑,但你在配置时还是要确保 seed 是固定的,不要用随机值。
实操心得:如果你的训练任务需要频繁重启(比如抢占式调度),建议把 num_samples 设大一些,并且在 checkpoint 里保存数据加载器的状态。这样重启后可以从断点继续,不需要从头开始。
6. 进阶话题:自定义数据集与扩展
6.1 接入自定义数据格式的完整流程
框架内置的数据格式不一定能满足所有需求。比如你有一些特殊格式的数据,或者需要在加载时做在线增强,就需要自定义数据集类。
自定义数据集的核心是实现三个方法:__len__、__getitem__和get_indexed_dataset。其中get_indexed_dataset是关键,它需要返回一个支持索引访问的对象。
from mindformers.dataset import BaseDataset class MyCustomDataset(BaseDataset): def __init__(self, data_path, seq_length, **kwargs): super().__init__(**kwargs) self.data_path = data_path self.seq_length = seq_length self._load_data() def _load_data(self): # 加载数据并构建索引 self.data = [] self.index = [] with open(self.data_path, 'r') as f: for line in f: tokens = self._process_line(line) self.index.append(len(self.data)) self.data.extend(tokens) def __len__(self): return len(self.index) - 1 def __getitem__(self, idx): start = self.index[idx] end = self.index[idx + 1] tokens = self.data[start:end] # padding 或截断到 seq_length return self._pad_or_truncate(tokens, self.seq_length)实现自定义数据集时要注意:__getitem__的返回值必须是固定长度的,通常是 seq_length。如果原始序列长度不足,需要 padding;如果超过,需要截断。padding 的值通常是 0 或者 tokenizer 的 pad_token_id。
6.2 多数据源权重动态调整的思路
固定权重在大多数场景下够用,但有些场景下你可能希望动态调整权重。比如训练初期多用通用数据,后期多用领域数据。这种课程学习的策略,可以通过自定义采样器来实现。
基本思路是:继承框架的采样器类,重写采样逻辑,根据当前训练步数动态计算权重。比如:
class DynamicWeightSampler: def __init__(self, datasets, initial_weights, final_weights, total_steps): self.datasets = datasets self.initial_weights = initial_weights self.final_weights = final_weights self.total_steps = total_steps self.current_step = 0 def get_weights(self): # 线性插值 ratio = min(self.current_step / self.total_steps, 1.0) weights = [ init * (1 - ratio) + final * ratio for init, final in zip(self.initial_weights, self.final_weights) ] return weights def sample(self, batch_size): weights = self.get_weights() # 按权重采样 ... self.current_step += 1这种动态权重策略在实际使用中需要谨慎,因为权重变化太剧烈可能导致训练不稳定。建议变化过程尽量平滑,并且总步数设置得足够长。
6.3 与 MindSpore 数据并行机制的配合
最后说一下 Blended DataLoader 和 MindSpore 数据并行机制的配合。在数据并行模式下,每个 rank 有独立的模型副本,但共享同一份数据的不同分片。
关键配置是dataset_strategy。在 MindSpore 里,可以通过mindspore.dataset.config.set_dataset_strategy来设置数据集的切分策略。对于 Blended DataLoader,通常设置为按 batch 维度切分:
import mindspore.dataset as ds ds.config.set_dataset_strategy( dataset_strategy="full_batch", num_shards=world_size, shard_id=rank_id )full_batch模式下,每个 rank 拿到完整的 batch,然后框架内部再做切分。这种方式的好处是数据加载逻辑简单,缺点是每个 rank 都要加载完整 batch 的数据,内存开销大。
另一种方式是data_parallel模式,每个 rank 只加载自己需要的那部分数据。这种方式内存效率高,但需要数据加载器支持按 rank 切分。
选择哪种方式取决于你的具体场景。如果 batch size 不大,内存充足,用full_batch更简单。如果 batch size 很大,内存紧张,用data_parallel更合适。
我在实际项目中的体会是,数据管道这块的工作量经常被低估。模型代码可能几天就写完了,但数据管道调通、调优可能要花一两周。尤其是分布式场景下,各种边界情况特别多。建议在项目初期就重视数据管道的设计和测试,不要等到训练跑不起来才回头排查。另外,数据格式尽量统一,不要混用多种格式,否则维护成本会成倍增加。