资讯动态

MindSpore Transformers 下 Megatron 风格数据集预处理实战

发布时间:2026/9/14 19:49:02 来源:尧图企业网站定制
前阵子帮团队迁移一个大模型训练链路数据端原本是一整套基于 Hugging Facedatasets的 JSONL 处理流程训练端却要从 PyTorch 切到 MindSpore。代码迁移本身倒还好数据格式先给我上了一课几百GB的 JSONL每轮 epoch 都要重新读取解析加载慢、内存涨幅吓人多卡并行时随机访问也费劲。后来我把预处理改成了 Megatron 风格的二进制格式——说白了就是一个 token 序列 bin 文件加一个索引 idx 文件训练吞吐直接上了一个台阶。这篇文章就把这套“MindSpore Transformers 环境下做 Megatron 风格数据集预处理”的思路和代码完整拆开讲一遍适合准备做 LLM 预训练、又纠结数据加载效率的团队参考。1. 为什么非要把数据整成 Megatron 格式1.1 Hugging Face 格式在大规模预训练里的短板用 Hugging Face 的datasets加载 JSONL 文本非常爽load_dataset(json, data_files...)就完事Arrow/parquet 底层帮你做了列式存储和缓存。但到了几百 TB token 级别的 LLM 预训练场景它有几个绕不开的短板。第一是随机访问效率。Hugging Face 的 Arrow 格式为了兼容各种数据类型和列式操作每次取一条样本要经过 schema 解析、列转行、内存拷贝。单卡实验没问题几十卡并行时每个 worker 都在频繁解析CPU 开销非常可观。第二是 shuffle 和重复 epoch 的体验。虽然datasets也支持shuffle和repeat但它的随机性控制、分片逻辑都是围绕“样本”而不是“token 流”做的。预训练模型需要的是在一个超长的 token 序列上不断滑窗切样本这种情况下把数据组织成“一条一条的文本”反而多了一层转换。第三是内存占用。JSONL 流式读取虽然内存小但每次迭代都要重读文本再 tokenize如果不想每轮都重新分词就得把 tokenize 结果缓存成 Arrow 或内存映射文件最终还是绕不开自定义二进制格式。1.2 Megatron 格式的核心设计token 流与文档索引分离NVIDIA Megatron-LM 里那套有名的数据预处理其实解决的就是上面三个问题。它的核心思路特别简单粗暴把整个数据集当成一个超长的 token 流来存。bin 文件所有文档的 token id 拼接成的一维数组内存映射mmap后可以零拷贝访问。数据在磁盘上就是连续排布访问任意位置就是一次指针偏移。idx 文件记录每个原始文档在 bin 里的偏移量。同时还会记录 vocab size、sequence length、总文档数等元信息。训练时查 idx 就能知道某个样本落在哪些文档范围内。因为 bin 是连续存储构造训练样本时就变成了纯数组切片操作随机取一个起点向前切seq_len个 token 就行。没有 JSON 解析没有字符串操作没有 Arrow 的列式转换性能和内存占用自然好看。1.3 理解文档边界为什么不能丢Megatron 格式里 idx 文件存储文档边界这件事很多从 Hugging Face 切过来的同学容易忽略。预训练时的样本是直接从 token 流里连续切出来的一个seq_len窗口完全可能横跨两个不同的原始文档。举个例子。文档 A 是 500 个 token文档 B 是 3000 个 token模型要 2048 的序列长度。切出来的某个样本可能是 A 的后 500 个 token 加上 B 的前 1548 个 token。如果不做任何标记模型就会把 B 的开头当成 A 的自然延续去学学到一堆压根不存在的“跨文档因果关系”。所以 Megatron 格式必须保留 document 边界训练时把样本内跨文档位置的 loss mask 置为 0或者通过 attention mask 禁止跨文档 attend。这也是后续在 MindSpore 里加载时最需要小心的地方。1.4 三种常见数据格式对比格式加载方式随机访问内存占用适合场景Hugging Face Arrow/parquet按行解析列式数据中等偏高微调、小规模训练、数据探索JSONL 流式读取逐行读文本再 tokenize低低数据清洗、格式转换Megatron bin idxmmap 零拷贝索引极高极低大规模预训练、跨框架切换我自己的体会是如果只是几百万条指令数据做微调HF 格式完全够用没必要折腾。但一旦进入 pre-training、连续多个 epoch、几百卡并行的阶段Megatron 风格的优势会指数级放大。2. 预处理全流程拆解从原始文本到二进制样本2.1 原始数据准备与清洗规则无论下游用什么格式原始数据清洗这一步不能省。我把常用规则列一下。统一编码为 UTF-8去掉 BOM 头。别小看这个BOM 混进 tokenizer 后经常会产生诡异的特殊 token。去除 HTML 标签、markdown 标记、重复空白字符。注意保留换行大多数代码模型对换行敏感。过滤过短文本比如少于 50 个 token和超长文本比如超过 10 万 token。超长文本可以直接切分但记录 document 边界时会多一层处理初期建议先切分或过滤。按文本哈希做全局去重尤其爬虫类数据重复率经常超 20%。清洗后统一转成 JSONL每行一个 JSON 对象至少包含{text: ...}字段。如果后续要做领域加权可以再加domain: code之类的字段预处理时顺手记录进 idx 扩展信息里。2.2 Tokenize 时的关键参数与 ID 类型选择清洗完后用transformers库加载 tokenizer。需要注意的点使用 fast tokenizer。AutoTokenizer.from_pretrained(..., use_fastTrue)尤其对英文和代码语料Rust 后端的提速非常明显。对 GPT 风格模型add_special_tokensFalse。如果你用的是 chat 类模型比如 Qwen 的 chat 版要特别注意它的 special token 设计这种情况建议参考对应模型的预训练数据处理脚本不要照搬。文本里如果已经有|endoftext|之类的分隔符要统一规则。常见做法是 tokenizer 正文不用特殊 token文档之间只在训练时用 attention mask 隔离。接下来是 ID 类型选择。这一步直接决定 bin 文件体积。判断逻辑很简单max_token_id tokenizer.vocab_size - 1 if max_token_id 65535: dtype np.uint16 else: dtype np.int32多数中文和英文 tokenizer 的 vocab size 都在 15 万以内用 uint16 就够。同样的 token 流uint16 比 int32 省一半磁盘和一半 IO。但切记如果某个 tokenizer 的 vocab 超过 65535uint16 存不下加载后会出现随机性的超大 id训练直接崩。这个坑我踩过后面问题排查部分会细讲。为什么省一半这么重要LLM 预训练数据集动辄几百 GB其实就是 token 数组。磁盘占用倒好说关键是训练时每个 step 要频繁读取这批数据IO 吞吐直接跟 bin 文件大小挂钩。省一半体积等于省一半传输时间对大规模训练是实打实的收益。2.3 拼接、切分与索引生成逻辑Tokenize 完成后每个文档变成一维 token id 数组。接下来的处理分三小步。第一步把每个文档的长度依次记录成sizes数组。然后计算累积偏移量doc_offsetsdoc_offsets[0] 0doc_offsets[i] doc_offsets[i-1] sizes[i-1]。这个偏移数组就是 idx 文件的核心。第二步把所有 token id 拼接成一个大数组。如果数据量不大可以直接np.concatenate如果数据量太大边读边写 bin 文件更好避免内存峰值。第三步确定样本切分方式。Megatron 常见的是无重叠连续切分从 token 流起点开始每隔seq_len切一个样本一直切到末尾。这样每个 token 在数据集中只出现一次训练效率最高。还有一种做法是带 stride 的滑动窗口切分让相邻样本有重叠能变相增加训练步数但样本之间有冗余一般不推荐在超大规模预训练里用。idx 文件的格式可以自己定义。我常用的一个版本magic4字节固定为MSPT versionint32 dtype_codeint321uint162int32 seq_lenint32 total_documentsint64 doc_offsetsint64数组长度 total_documents1 sizesint64数组长度 total_documents这里把 doc_offsets 和 sizes 都保存一份是为了在 MindSpore 训练时快速判断样本是否跨文档。只存 offsets 也能推出 sizes但提前存好可以省去每次加载时的减法运算。3. 实操写一个自己的 Megatron 风格预处理脚本3.1 从 JSONL 到 bin idx 的完整代码下面这段代码可以直接用。假设输入是一个 JSONL 文件每行{text: ...}输出是output_prefix.bin和output_prefix.idx。import argparse import json import struct import numpy as np from pathlib import Path from transformers import AutoTokenizer def read_documents(input_path: str): docs [] with open(input_path, r, encodingutf-8) as f: for line in f: line line.strip() if not line: continue obj json.loads(line) text obj.get(text, ) if len(text) 20: continue docs.append(text) return docs def tokenize_documents(docs, tokenizer): sizes [] all_tokens [] for text in docs: ids tokenizer.encode(text, add_special_tokensFalse) if len(ids) 0: continue sizes.append(len(ids)) all_tokens.extend(ids) return all_tokens, sizes def write_bin_idx(output_prefix: str, all_tokens, sizes, seq_len, dtype): tokens_np np.asarray(all_tokens, dtypedtype) bin_path f{output_prefix}.bin idx_path f{output_prefix}.idx # 写 bin 文件 tokens_np.tofile(bin_path) # 写 idx 文件 offsets np.zeros(len(sizes) 1, dtypenp.int64) for i, s in enumerate(sizes): offsets[i 1] offsets[i] s dtype_code 1 if dtype np.uint16 else 2 total_docs len(sizes) with open(idx_path, wb) as f: f.write(bMSPT) f.write(struct.pack(i, 1)) # version f.write(struct.pack(i, dtype_code)) # dtype_code f.write(struct.pack(i, seq_len)) # seq_len f.write(struct.pack(q, total_docs)) # total_documents offsets.tofile(f) np.asarray(sizes, dtypenp.int64).tofile(f) print(ftotal tokens: {len(tokens_np)}) print(ftotal docs: {total_docs}) print(festimated samples: {len(tokens_np) // seq_len}) def main(): parser argparse.ArgumentParser() parser.add_argument(--input_path, typestr, requiredTrue) parser.add_argument(--tokenizer, typestr, requiredTrue) parser.add_argument(--output_prefix, typestr, requiredTrue) parser.add_argument(--seq_len, typeint, default2048) parser.add_argument(--num_workers, typeint, default8) args parser.parse_args() tokenizer AutoTokenizer.from_pretrained(args.tokenizer, use_fastTrue) max_token_id tokenizer.vocab_size - 1 dtype np.uint16 if max_token_id 65535 else np.int32 docs read_documents(args.input_path) all_tokens, sizes tokenize_documents(docs, tokenizer) write_bin_idx(args.output_prefix, all_tokens, sizes, args.seq_len, dtype) if __name__ __main__: main()这段代码刻意写得简单直白方便理解逻辑。实际处理超大语料时all_tokens列表可能占掉大量内存建议改成“边 tokenize 边写入 bin 临时文件”最后再扫描临时文件生成 offsets。也可以直接用多进程加速 tokenize比如用multiprocessing.Pool把文档分片处理。3.2 性能建议与运行示例运行命令示例python preprocess_megatron_style.py \ --input_path /data/corpus.jsonl \ --tokenizer Qwen/Qwen2.5-7B \ --output_prefix /data/output/qwen_7b \ --seq_len 4096 \ --num_workers 16实测下来一个中等规模的 tokenizer比如 vocab 15 万左右单进程 tokenize 速度大约在每秒 10 万到 20 万 token开启 fast tokenizer 后会更快。瓶颈通常不在分词器而在 JSON 解析和磁盘 IO。如果原始数据是 gzip 压缩的建议先解压再处理否则 CPU 全耗在解压上得不偿失。对于大规模语料我特别推荐分 shard 处理。把整个语料按 5GB 一个文件切分每个 shard 单独跑一遍预处理生成 bin 和 idx之后训练时再在 sample 层面做多文件联合随机采样。这样既方便并行处理也方便后续做数据混入比例的调整。3.3 生成后的验证工具预处理完一定要做验证不要直接拿去训练。我常用的验证脚本import numpy as np def load_idx(idx_path): with open(idx_path, rb) as f: magic f.read(4) assert magic bMSPT, invalid index file version np.fromfile(f, dtypenp.int32, count1)[0] dtype_code np.fromfile(f, dtypenp.int32, count1)[0] seq_len np.fromfile(f, dtypenp.int32, count1)[0] total_docs np.fromfile(f, dtypenp.int64, count1)[0] offsets np.fromfile(f, dtypenp.int64, counttotal_docs 1) sizes np.fromfile(f, dtypenp.int64, counttotal_docs) return { version: version, dtype_code: dtype_code, seq_len: seq_len, total_docs: total_docs, offsets: offsets, sizes: sizes, } def validate(bin_path, idx_path): idx load_idx(idx_path) dtype np.uint16 if idx[dtype_code] 1 else np.int32 tokens np.memmap(bin_path, dtypedtype, moder) total_tokens tokens.shape[0] assert total_tokens idx[offsets][-1], token count mismatch print(total tokens check passed) # 随机抽几个起点验证能否正确切片 for start in [0, 100, 100000, total_tokens // 2]: if start 2048 total_tokens: continue sample tokens[start: start 2048] assert sample.shape[0] 2048 print(fsample from {start} ok, min_id{sample.min()}, max_id{sample.max()})如果max_id接近甚至超过 65535而当时用的是 uint16这里就能立刻发现不用等到训练崩溃。4. 在 MindSpore 训练中接入 Megatron 风格数据4.1 用 GeneratorDataset 包装 bin idxMindSpore Transformers也就是 mindformers 库本身支持多种数据源但直接读自定义二进制文件最灵活的方式还是继承数据源后传给GeneratorDataset。我封装了一个轻量类import numpy as np import mindspore as ms from mindspore.dataset import GeneratorDataset class MegatronStyleDataset: def __init__(self, bin_path, idx_path, seq_len, sample_offsetsNone): self.seq_len seq_len idx self._load_idx(idx_path) self.dtype np.uint16 if idx[dtype_code] 1 else np.int32 self.tokens np.memmap(bin_path, dtypeself.dtype, moder) self.doc_offsets idx[offsets] self.doc_sizes idx[sizes] self.total_tokens self.tokens.shape[0] self.valid_samples self.total_tokens // self.seq_len if sample_offsets is None: # 无重叠切分起点为 i * seq_len self.sample_offsets np.arange(0, self.valid_samples * self.seq_len, self.seq_len) else: self.sample_offsets sample_offsets def _load_idx(self, idx_path): with open(idx_path, rb) as f: assert f.read(4) bMSPT, invalid idx file version np.fromfile(f, dtypenp.int32, count1)[0] dtype_code np.fromfile(f, dtypenp.int32, count1)[0] seq_len np.fromfile(f, dtypenp.int32, count1)[0] total_docs np.fromfile(f, dtypenp.int64, count1)[0] offsets np.fromfile(f, dtypenp.int64, counttotal_docs 1) sizes np.fromfile(f, dtypenp.int64, counttotal_docs) return { version: version, dtype_code: dtype_code, seq_len: seq_len, total_docs: total_docs, offsets: offsets, sizes: sizes, } def __len__(self): return len(self.sample_offsets) def __getitem__(self, index): start int(self.sample_offsets[index]) end start self.seq_len input_ids np.array(self.tokens[start:end], dtypenp.int32) # 生成 loss_mask跨文档的位置置 0 loss_mask np.ones(self.seq_len, dtypenp.float32) doc_idx np.searchsorted(self.doc_offsets, start, sideright) - 1 if end self.doc_offsets[doc_idx 1]: pass # 整个样本都在一个文档内 else: # 从 doc_idx 开始逐个文档判断 pos start mask_accum [] current_doc doc_idx while pos end: doc_end self.doc_offsets[current_doc 1] piece_end min(end, doc_end) mask_accum.append((pos - start, piece_end - start)) pos piece_end current_doc 1 for piece_start, piece_end in mask_accum: # 跨文档的第二个及之后的片段loss_mask 置 0 if piece_start 0: loss_mask[piece_start:piece_end] 0.0 labels input_ids.copy() return input_ids, labels, loss_mask def build_megatron_dataset(bin_path, idx_path, seq_len, batch_size, shuffleTrue): dataset MegatronStyleDataset(bin_path, idx_path, seq_len) ds GeneratorDataset( sourcedataset, column_names[input_ids, labels, loss_mask], shuffleshuffle, ) ds ds.batch(batch_size) return ds这里__getitem__返回了三样东西input_ids、labels、loss_mask。labels在多数自回归模型里就是input_ids的右移这一步可以直接在模型内部处理我返回一份是为了灵活性。loss_mask是关键跨文档的位置被置为 0loss 计算时这些 token 不贡献梯度模型就不会学到虚假的跨文档依赖。4.2 样本切分与跨文档 Mask 逻辑的细节上面代码里loss_mask的生成逻辑值得单独说两句。我先用searchsorted找到起点落在哪个文档然后从那个文档开始循环把样本按文档边界切成若干段。第一段保留 loss后面所有段都置 0。为什么第一段保留因为第一段是从这个文档中间开始的它前面的上下文确实属于同一个文档模型可以正确建模。后面几段属于新文档虽然不能attend到前一个文档但新文档内部仍然可以正常建模。更严格的实现还应该生成 attention mask让整个样本内只有同一文档的 token 能互相 attend。但在大多数主流模型实现里用loss_mask置 0 已经足以避免跨文档学习。如果你用的是 mindformers 的GPT2LMHeadModel这类封装好的模型它通常把labels和loss_mask都传给 loss 函数置 0 的位置在 cross entropy 里自动被忽略效果一样。4.3 与 mindformers 训练脚本的对接在 mindformers 里你可以在现有训练流程中替换数据加载部分。比如训练脚本里原本用build_dataset读原始文本现在可以直接用我上面写的build_megatron_dataset。以 mindformers 的配置文件为例你只需要把dataset那一块的type指向自定义的逻辑。碰到直接用 Python API 写训练循环的情况更简单直接在model.train或自定义Trainer里传入build_megatron_dataset返回的ds就行。不需要改模型结构不需要改优化器只要保证喂进去的input_ids形状是(batch, seq_len)loss_mask形状也是(batch, seq_len)。有一个性能细节要提醒GeneratorDataset默认的num_parallel_workers是 1多卡训练时建议调大比如 8 或者 16。同时如果数据集特别大预生成sample_offsets比在__getitem__里现场随机生成要稳定因为 mindspore 的多进程 worker 各自持有不同的随机状态容易出现重复采样。我的做法是在主进程一次性生成 offset 数组然后shuffleTrue交给GeneratorDataset处理每个 epoch 内部会重新 shuffle。5. 常见问题与排查方法5.1 高频问题速查表我在实际迁移中遇到过不少问题整理成一张表方便排查。现象原因解决方式加载 bin 后 token id 出现超大值或负数使用了 uint16 但 vocab_size 超过 65535检查 tokenizer vocab改用 uint32/int32样本开头频繁出现[UNK]或\x00原始文本里有非法字符或忘记add_special_tokensFalse清洗数据编码时指定errorsignore确保参数正确训练第一个 step loss 就变成 NaNloss_mask 全为 0 或 attention mask 异常验证 loss_mask 是否只对跨文档片段置 0检查是否有 padding token 参与了 loss 计算多卡训练时不同卡数据高度重复__getitem__内部随机生成 offsetworker 间随机状态冲突预先在主进程生成全部 offset 数组交给 Dataset 做 shufflemmap 加载后读取内容与源文件不符dtype 不匹配或 bin 文件路径错误用np.fromfile对照原始 token 验证检查 endian训练吞吐提升不明显GeneratorDataset的num_parallel_workers太小调大 worker 数并开启数据管道的 prefetch5.2 两个值得展开的实战案例第一个是 uint16 溢出。有次我用一个多语种 tokenizer 训练vocab size 正好卡在 65536 附近。预处理脚本里判断条件是max_token_id 65535按道理该走到 int32 分支。但那段时间我为了省显存手滑在另一个配置里强制指定了uint16。训练跑了一万步都没事loss 正常下降直到某次采样到 vocab 尾部的 token模型直接崩掉。排查了很久才发现是 token id 被截断。建议在预处理脚本里加一个断言assert tokens_np.max() 65535 or dtype ! np.uint16。第二个是跨文档 mask 没生效。之前用简单实现只把loss_mask在跨文档片段置 0但忘了处理“多个文档边界出现在同一个样本”的情况。一个样本横跨三个文档时第二段和第三段都要置 0我只处理了第二段。结果模型在部分样本里还是学到了跨文档依赖下游任务指标略降。后来改成循环遍历逐个文档片段判断问题才解决。5.3 关于数据加载性能的几个深度调优点除了表格里的问题正式跑大规模训练前还可以检查这几处样本起始偏移的分布。如果所有 offset 都是i * seq_len样本之间完全无重叠训练速度快但每个 token 只被看到一次。有些实验想增加样本数量会故意在每个 epoch 打乱 offset。建议先固定无重叠之后再扩展。bin 文件在机械硬盘上的顺序读性能远好于随机读。如果用 HDD 存数据尽量让 Dataloader 按顺序批量读再用内存 shuffle。SSD 上可以直接随机读但多 worker 并发时也要控制 IO 队列深度。如果 MindSpore 版本支持mindspore.dataset的config.set_prefetch_size可以适当调大 prefetch减少训练 step 间等待数据的时间。6. 从零到一接入的完整建议如果你正要给 MindSpore 上的 LLM 训练切换数据格式我建议按这个顺序来先拿一个小数据集比如 1 万条文本几十 MB跑通完整的“JSONL - tokenize - binidx - GeneratorDataset - 训练脚本”确认 loss 曲线正常再扩展到大语料。不要一上来就全量转换数据格式的问题在小数据上几分钟就能暴露全量了可能折腾一天。扩展到大语料时优先考虑多 shard。每个 shard 独立生成 bin 和 idx训练时多个 shard 交替读取。这样能做数据混入比例控制比如代码数据 70%、文本数据 30%同时每个 shard 内部的随机访问依然高效。我在实际使用中还有一个习惯把 tokenizer 的名字、vocab 大小、dtype、seq_len 全部写进 idx 文件头。虽然占不了几个字节但防止半年后自己都忘了这份数据是用哪个 tokenizer 生成的。版本信息是救命稻草尤其是团队协作时数据文件比代码更不容易追溯来源。这套方案跑通之后最直观的感受是数据加载从“训练瓶颈”变成了“完全无感”。同样的硬件不吃内存不频繁解析多卡扩展也平稳。个人经验是凡是训练数据超过 100GB 的 LLM 任务都值得把预处理切到 Megatron 风格。如果你也在 MindSpore 上折腾 LLM照这个思路做一遍应该能少踩不少坑。

读完文章,也想定制专属网站?

尧图设计师 24 小时内与您沟通定制方案

免费获取报价