1. 大模型训练里最容易被低估的环节数据管道做过大模型预训练的人都有一个共识模型结构决定上限数据质量决定下限而数据管道的效率直接决定你多久能摸到这个下限。我见过太多团队在模型并行、算子优化上砸了几周时间结果训练吞吐上不去最后定位下来问题出在 DataLoader 上——GPU 利用率长期在 60% 上下晃荡算力全耗在等数据了。MindSpore Transformers 这套框架里的 Blended Megatron DataLoader就是专门解决这个问题的。它做的事情说起来不复杂把多个不同来源、不同格式的数据集按指定权重混合成一个统一的训练数据流同时保证在分布式训练场景下每个卡拿到的数据既不重复也不遗漏还要让数据预取和计算重叠起来把 I/O 等待藏到计算背后。这套机制适合谁看如果你正在用 MindSpore 做 LLM 预训练或者微调需要自己准备和处理训练数据或者你发现训练时 step time 波动大、GPU 利用率上不去那这篇文章里的内容应该能帮到你。即便你用的是别的框架Megatron 系的数据加载思路也是通用的理解了这个设计换个框架照样能迁移。我下面会从整体设计思路开始拆然后逐层深入到配置细节、实操步骤、参数计算最后把我踩过的坑和排查经验整理出来。内容偏实操代码和配置都会给到可以直接参考的版本。2. 整体设计思路为什么要做 Blended DataLoader2.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, dtypenp.int64) for i, tokens in enumerate(token_ids_list): offsets[i 1] offsets[i] len(tokens) # 索引文件格式前两个值存维度信息后续存偏移 header np.array([len(token_ids_list), 1], dtypenp.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 为 Wglobal_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 DataLoader4.1 数据准备与预处理脚本假设我们有三份原始数据web_data.jsonl、code_data.jsonl、book_data.jsonl。每行是一个 json 对象包含 text 字段。我们需要把它们转换成 Megatron 二进制格式。第一步是 tokenize。这里用 MindSpore Transformers 自带的 tokenizerfrom mindformers import AutoTokenizer tokenizer AutoTokenizer.from_pretrained(llama2_7b) def tokenize_file(input_path, output_path, max_seq_len2048): 将 jsonl 文件 tokenize 并写入二进制文件 all_tokens [] with open(input_path, r, encodingutf-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], dtypenp.int32) offsets np.zeros(len(lengths) 1, dtypenp.int64) offsets[1:] np.cumsum(lengths) # 写入索引文件 with open(f{output_prefix}.idx, wb) as f: # 头部样本数、维度 f.write(np.array([len(lengths), 1], dtypenp.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_lengthRuntimeError: batch size not divisibleglobal_batch_size 不能被 world_size 整除调整 global_batch_size 或 world_size训练 loss 为 NaNtoken 值超出 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_strategyfull_batch, num_shardsworld_size, shard_idrank_id )full_batch模式下每个 rank 拿到完整的 batch然后框架内部再做切分。这种方式的好处是数据加载逻辑简单缺点是每个 rank 都要加载完整 batch 的数据内存开销大。另一种方式是data_parallel模式每个 rank 只加载自己需要的那部分数据。这种方式内存效率高但需要数据加载器支持按 rank 切分。选择哪种方式取决于你的具体场景。如果 batch size 不大内存充足用full_batch更简单。如果 batch size 很大内存紧张用data_parallel更合适。我在实际项目中的体会是数据管道这块的工作量经常被低估。模型代码可能几天就写完了但数据管道调通、调优可能要花一两周。尤其是分布式场景下各种边界情况特别多。建议在项目初期就重视数据管道的设计和测试不要等到训练跑不起来才回头排查。另外数据格式尽量统一不要混用多种格式否则维护成本会成倍增加。
阅读完成 · 觉得有帮助?