用MindSpore数据管道重构大模型数据预处理流程
发布时间:2026/9/30 13:12:05来源:尧图网络
1. 为什么我放弃了先预处理再训练的老路子1.1 踩坑背景数据脚本和训练脚本脱节之后最早接触昇思 MindSpore 的时候我和很多朋友一样习惯用传统三件套来做大模型的数据准备先写一个独立的 Python 脚本用 pandas 或者 numpy 把文本清洗、截断、构建 input_ids跑完导出成 npy 或者 json再在训练脚本里用另一个函数读进来转成 Tensor最后才丢给模型。这套流程在小数据集上看着没毛病但一旦涉及几十 GB 甚至上 TB 的原始语料问题就一个接一个冒出来。第一个问题是最典型的清洗规则改了一版结果只改了预处理脚本忘了同步训练脚本里的读取逻辑。两边对字段的假设不一样模型在训练到第 2000 步的时候突然崩掉报一个 shape mismatch查了半天才发现是某条样本在清洗阶段被截断成了不同长度。第二个问题是内存。很多 NLP 预处理的常规操作比如分词、padding、构建注意力掩码如果全部在训练循环外面提前做完相当于把全量数据集在内存里复制了好几份。我实测过一个 30GB 的文本数据集用 pandas 分批读完再合并内存峰值直接冲到 70GB 以上机器直接卡死。后来我才真正理解了一件事在大模型训练场景下数据预处理不应该是一个独立于训练流程之外的阶段而应该是数据集本身的一部分。昇思 MindSpore 的mindspore.dataset模块核心价值就是把数据变换、批次构造、分布式切分这些能力直接内嵌到训练管线里让框架自己去调度而不是靠人肉写循环。1.2 换一种思路让数据变换成为计算图的一部分MindSpore 里的GeneratorDataset可以把你的自定义数据源包装成一个 Dataset 对象然后通过.map()、.batch()、.shuffle()、.repeat()这些算子组合成一条完整的数据流水线。这听起来像是 PyTorch 的 DataLoader 或者 TensorFlow 的 tf.data但实际用起来有个很明显的差别MindSpore 的 Dataset 算子在底层是跟着计算图一起优化的它知道哪些变换可以合并哪些可以异步执行哪些可以在昇腾 NPU 上直接做。我在一次视觉大模型的数据准备中试过把 50 万张图片的随机裁剪、归一化、色彩抖动全部放进 Dataset pipeline 里而不是预处理成离线文件。训练时 NPU 的利用率比之前用离线预处理数据高了不少因为图像变换被分发到了设备端的多核处理器上CPU 和 NPU 的流水线重叠得更好。这个收益在大模型训练这种数据吞吐要求高于计算瓶颈的场景里非常可观。所以这篇文章的核心内容就是一套基于mindspore.dataset的数据变换与预处理方案。我会从管道设计思路讲起再到常用的算子组合顺序、文本和多模态数据的处理策略、分布式训练下的切分细节最后聊一聊我在 VS Code 里的调试方法和踩过的坑。无论你是在昇腾 NPU 上跑大模型还是想在 GPU 环境下把数据预处理做得更规整这套思路都值得参考。2. 数据变换的底层思路先画管道再写代码2.1 和 NumPy/pandas 预处理流程的本质差异刚开始切换到我推荐用mindspore.dataset管线时同组同事最常问我的一句话是我在预处理脚本里做完再存下来和你在 Dataset 里现做有什么区别区别可以从三个维度来看。第一是执行时机。离线预处理是在训练开始前一次性执行完的结果以文件形式落盘Dataset 管道则是在每个 epoch 内按需执行并且可以和训练计算重叠。这一点决定了你的数据流水线能不能跑满加速卡。第二是调度粒度。离线脚本是批处理apply()一个函数就能对整列数据做操作而 Dataset 管道的粒度是样本后变换每个算子都面向单个样本或小批次执行这样你可以针对单条坏数据做精细处理而不用因为一条脏数据重跑全量脚本。第三是随机性控制。离线预处理如果保存的是已经做好的随机增强结果每个 epoch 看到的数据其实是同一份Dataset 管道则能在每次访问时重新执行随机变换天然支持每次 epoch 看到不同的增强结果这在数据量有限的大模型微调场景下特别重要。我画过一张简单的对比表把常用操作在两种思路下的位置列出来操作类型离线预处理脚本mindspore.dataset 管道文本清洗一次性执行存盘map中的PyFunc按样本执行随机裁剪预先裁剪增强效果固化map或batch内动态执行Padding定长 padding浪费存储batch阶段动态 padding节省存储数据切分手动按比例切分文件split或分布式采样器自动切分Shuffle只能文件级乱序shuffle算子支持缓冲区级乱序从表格里能看出离线脚本更适合一批次全量处理、结果固定的场景而 Dataset 管道更适合数据量大、需要动态增强、训练过程长的场景。大模型训练显然属于后者。2.2 一个最小可用示例原始文本文件到 Dataset 对象在mindspore.dataset里最常见的起点是GeneratorDataset。它接收一个 Python 迭代器或者生成器函数把取出来的每一条数据转换成 MindSpore 能识别的张量。下面是一个最小可用的示例import mindspore as ms import mindspore.dataset as ds def text_samples(): # 模拟从磁盘读取原始文本 for line in open(corpus.txt, r, encodingutf-8): if line.strip(): yield {text: line.strip(), label: 1} dataset ds.GeneratorDataset( sourcetext_samples(), column_names[text, label], shuffleFalse ) # 查看前两条样本 for data in dataset.take(2): print(data)column_names这一步很重要它相当于给每条样本的字段起了个名字后面的map算子才能通过字段名指定要对哪一列做变换。比如要对text这列做分词def tokenize(text): # 这里可以是任意的 tokenizer 逻辑 return list(text) dataset dataset.map(operations[tokenize], input_columns[text])这里input_columns如果不写默认会对全部字段做操作如果只写一个字段其他字段保持不变。很多人刚上手时容易在input_columns和output_columns上犯迷糊搞不清楚变换之后新列叫什么名字。我的经验是output_columns只在input_columns数量与output_columns数量不一致时才需要显式设置平时保持默认就好。实际项目里构建 Dataset 对象之前最好再加一层原始数据校验。比如检查文件是否存在、列名是否齐全。数据源如果是从多个文件拼起来的更建议用glob先做一次文件列表排序避免不同进程读取顺序不一致。3. 大模型场景下常用算子组合与先后顺序3.1 为什么顺序这么重要shuffle、batch、map、repeat 的配合逻辑很多人在用mindspore.dataset的时候觉得算子顺序无所谓反正都是变换。但真实训练场景里顺序错一个结果可能完全不同而且很难排查。核心顺序法则是先 map 后 batch先 shuffle 后 mapbatch 之后再考虑 repeat。下面逐一解释。先 map 后 batch是因为绝大多数逐样本变换文本截断、图像缩放、分词在单条数据上执行成本更低。如果先 batch那么 map 算子就要对整个 batch 内的数据做操作不仅要处理嵌套结构还可能造成内存峰值上升。MindSpore 的map算子可以自动处理batch之后的数据结构但计算效率并不会更好。先 shuffle 后 map是为了防止预处理结果被永久绑定到原始顺序。比如你做随机遮挡如果先 shuffle 再 map每个 epoch 里样本顺序不同遮挡结果也会随之变化如果先 map 再 shuffle增强结果被固化在缓存里shuffle 只改变排列不改变增强结果。从数据多样性角度讲先 shuffle 后 map 更有利。batch 之后再 repeat这是一个非常容易被忽视的细节。repeat()并不是简单地重复迭代数据集它会把整个 pipeline包括已经执行的 map 和 shuffle重复多次。如果 shuffle 放在 repeat 之前MindSpore 会为 shuffle 算子维护一个较大的缓冲区每次 repeat 都能重新洗牌如果 repeat 放在 shuffle 之前相当于把同一个 shuffle 后的序列重复多遍数据顺序完全一样等于没洗牌。我整理了一个常见的顺序模板适合大模型预训练或者微调dataset ds.GeneratorDataset(source, column_names[text, label]) dataset dataset.shuffle(buffer_size10000) dataset dataset.map(operations[clean_text], input_columns[text]) dataset dataset.map(operations[tokenize_and_pad], input_columns[text]) dataset dataset.batch(batch_size32, drop_remainderTrue) dataset dataset.repeat(epochs)这里drop_remainderTrue是为了保证最后一个 batch 不会出现比batch_size小的样本免得模型在 unsqueeze 或者计算 attention mask 时出问题。如果你用的是动态 shape 的网络可以设为False但内部最好加上长度校验。3.2 文本与多模态输入的 mask、padding 与截断策略大模型的输入一般是以 token id 为主甚至还会用到attention_mask、token_type_ids、labels多模态场景还得有pixel_values。在mindspore.dataset中做这些操作最优雅的方式是把 tokenizer 封装成一个 map 函数。def process_text(row): text row[text] encoded tokenizer( text, max_length512, paddingmax_length, truncationTrue, return_tensorsnp ) return encoded[input_ids][0], encoded[attention_mask][0] dataset dataset.map( operations[process_text], input_columns[text], output_columns[input_ids, attention_mask] )需要注意的一个细节是tokenizer如果在map内部实例化每次执行变换时都会重新初始化导致性能损耗。所以在构建 Dataset 之前先在外面把tokenizer实例化一次map 函数内部直接调用它不重复加载模型。这一点和GeneratorDataset的 lazy 执行机制是兼容的因为 map 函数每次被调用时持有的是同一个闭包引用。对于 padding 和截断我强烈建议在map阶段就做好而不是等到batch阶段再动态处理。原因很简单大模型训练里的 attention mask 依赖于每个样本的实际长度如果在 batch 阶段统一 padding你就需要额外记录一个lengths列逻辑复杂得多。而在map阶段直接 padding 到固定长度后面所有算子都不用再关心长度变化。多模态输入的应对思路类似但多了一步类型转换。图像数据通过Decode解码成 uint8 数组再通过归一化和标准化算子转成 float32。从我的实践看图像 map 的算子数量和链式顺序对训练吞吐影响很大。建议把解码和缩放这种重 IO 算子放在最前面标准化和归一化算子放在后面中间不需要穿插其他字段的变换。3.3 性能取舍batch_size、num_parallel_workers 与 cache 的平衡大模型的数据管线最怕两个瓶颈一是 CPU 预处理太慢导致加速卡一直空转二是内存占用太大导致 OOM。MindSpore 的map算子提供了一个num_parallel_workers参数可以并行执行变换。这个参数该设多大取决于你机器的 CPU 核心数和单个样本变换的耗时。我自己的经验公式是num_parallel_workers设置为物理核心数的一半到三分之一比较稳妥。比如 32 核机器设为 8 到 12 比较合理。设得过大线程切换的开销会超过并行收益设得过小流水线又会频繁被打断。另一个值得关注的是cache算子。mindspore.dataset从 1.6 版本开始支持cache可以把 map 之后的中间结果缓存到内存或磁盘。如果数据集的某个部分比如验证集会被反复读取加缓存的效果立竿见影。cache ds.DatasetCache(session_id12345, size6*1024*1024*1024) dataset dataset.map(operations[process_text], ...).cache(cache)注意cache之后不能再接会改变样本数量的算子比如repeat或者是随机性很强的 map否则缓存的结果会不一致。这一点我踩过好几次坑给shuffle之前的数据加了缓存导致每个 epoch 读到的顺序完全一致训练效果直接下降。如果你用的是昇腾 NPU还有一个额外的小技巧数据管道的末端可以用to_device()或者直接交给model.train()来处理让 MindSpore 自动把数据搬运到设备端。不要自己在脚本里加as_encoder之类的操作多一层拷贝就是多一倍延迟。4. 数据清洗与脏数据兜底卫生习惯比模型结构更影响效果4.1 跳过错乱样本与空实例真实语料永远比想象中脏。我处理过大模型微调用数据里面经常出现空行、重复样本、长度过短/过长的样本甚至还有编码异常导致乱码的情况。如果这些脏数据直接进入训练流程模型可能在某个 batch 里被迫学习无意义的 attention 模式损失曲线出现诡异的尖刺。在map函数里做清洗最常见的方式是返回一个标记值再配合 filter 算子过滤。但 filter 算子会改变数据集的长度使用时要格外小心尤其是在有 cache 或 repeat 的场景下。另一个更稳妥的方案是在源头清洗在GeneratorDataset的生成器内部就把条件过滤掉。def text_samples(): seen set() for line in open(corpus.txt, r, encodingutf-8): line line.strip() if not line: continue if len(line) 5: continue if line in seen: continue seen.add(line) yield {text: line, label: 1}这个方法虽然原始但它保证了下游所有算子都看不到脏样本也不会有 filter 之后样本数变化带来的复杂边界情况。唯一的代价是seen集合会占用内存。我一般改为在源文件级别去重或者用sqlite临时表做 KEY 去重避免seen集合过于膨胀。对于长度过短的判定你最好先跑一次统计看真实的长度分布。不同的任务对短文本容忍度完全不同指令微调里 5~10 个 token 的样本可能反而是高质量数据但预训练语言模型里3 个 token 的碎句子基本没有学习价值。所以先统计再定阈值别拍脑袋。4.2 分布式训练下的每卡数据划分与 seed 的设定大模型训练基本逃不掉分布式。MindSpore 的 Dataset 在分布式下有两种切分方式一种是使用shard参数另一种是使用DistributedSampler。我用得比较多的是DistributedSampler它支持shuffle和seed两个关键参数。其中seed一定要固定否则每个 epoch 的洗牌顺序在各卡上不一致会引发数据交叉重复的问题。sampler ds.DistributedSampler( num_shardsrank_size, shard_idrank_id, shuffleTrue, seed42 ) dataset ds.GeneratorDataset( sourcetext_samples(), column_names[text, label], samplersampler )这里还有个容易被忽略的细节GeneratorDataset内部如果自带shuffle再和DistributedSampler的 shuffle 叠加会产生双重洗牌顺序变得不可控。我的建议是分布式训练时优先使用 sampler 控制洗牌Dataset 自带的shuffle参数设为 False。这样可以保证逻辑清晰各卡之间的样本不会串。还有一点在多机多卡场景下seed之后如果再修改全局mindspore.set_seed(42)那么所有采样器和随机操作的种子会相对统一方便复现。但要注意别在set_seed之后再调用一次set_seed否则后续随机序列会发生偏移。4.3 脏数据兜底宁可丢弃不可带病训练在map阶段如果一条样本的 tokenizer 返回结果异常比如所有 token 都是[UNK]或者 input_ids 长度为 0最简单的兜底策略不是修正它而是在源头丢弃。因为一条异常样本修正后再进入训练往往还带着隐蔽的格式问题后续处理成本更高。我在实际项目中是这么做的给process_text增加一层异常捕获异常样本直接返回一个特殊标记然后在生成器内部过滤掉。def process_text(row): try: text row[text] encoded tokenizer(text, max_length512, truncationTrue) if len(encoded[input_ids]) 0: return None return encoded[input_ids] except Exception: return None def safe_map(row): result process_text(row) if result is None: return None return result, row[label]不过要提醒一下如果map算子返回None后续的batch算子很可能会报错因为 None 不能被转换成张量。所以我通常不用这种方式而是在GeneratorDataset的生成器里处理所有异常逻辑保证进入 Dataset 的每条样本都是干净且结构完整的。这样后续的 map 就可以安心做业务变换不用再考虑兜底。5. 调试、验证与工具链整合5.1 在 VS Code 中使用 MindSpore 内核快速验证数据形状大模型开发里数据集管的调试往往比模型调试更费时间。我自己就是在 VS Code 里用 MindSpore 内核也就是 Jupyter 模式的 MindSpore 内核来做 pipeline 验证的。典型操作是先在一个 cell 里构建 Dataset 对象再用take(1)取一条数据打印出来看结构。这样可以快速确认字段名、类型、shape 是否符合预期而不用启动一个完整的训练任务。import mindspore.dataset as ds # 构建期间可以先不加载全量数据 dataset ds.TextFileDataset(corpus_sample.txt, shuffleFalse) sample dataset.take(1).create_dict_iterator().get_next() print(sample.keys()) print(sample[text].shape)这种方式对排查 shape mismatch 特别有用模型报错时说某个维度不对你先用相同参数取一条真实数据看看它的 shape 到底是什么基本能定位是map输出维度不一致还是batch的pad_info配置有误。不过在 VS Code 里跑 MindSpore 内核时有一个小坑如果你在 Notebook 里定义了GeneratorDataset然后重复执行同一个 cell生成器函数会重新被调用但数据集对象如果存的是同一个引用可能拿到的是迭代器状态的延续而不是新的数据。遇到这种情况直接重启内核再来一次就好不必纠结。5.2 用 create_dict_iterator 检查管道输出质量想验证整条数据管道的输出是否符合预期最直接的方式是create_dict_iterator。它在保留 dataset 原始结构的同时把每条样本变成一个 dict方便按字段查看。我的例行检查清单是这样检查input_ids长度是否等于attention_mask长度。检查attention_mask里 0 和 1 的分布是否符合 padding 预期特别是尾部 padding 后 mask 应为 0。检查labels列是否存在、shape 是否与input_ids一致。随机抽 20 条样本打印出来人工看一眼语义是否仍然可读。这一步基本能筛掉八成以上的管道问题。有时候问题出在 tokenizer 本身比如某些说话风格里带了很多特殊符号tokenizer 把它全拆成了[UNK]模型根本学不到东西。打印出来肉眼能看到的话就能尽早发现。for i, data in enumerate(dataset.create_dict_iterator(num_epochs1)): if i 20: print(fsample {i}: input_ids_length{len(data[input_ids])} fmask_sum{sum(data[attention_mask])})5.3 从数据管道导出上采样批次用于模型预热有一个比较少见但很实用的操作当你想在正式训练开始前先用一小批真实数据跑一遍模型的前向和反向验证网络结构没问题可以直接从 Dataset 里采样一个批次转成 numpy 后喂给模型。这样既不用写单独的测试脚本也不用去读一个离线文件。batch next(dataset.create_tuple_iterator(num_epochs1)) model_input, label batch output model(model_input) loss loss_fn(output, label) loss.backward()这个做法在 Lightweight 验证阶段非常管用。我通常把这段代码放在训练脚本的入口通过一个--smoke-test参数来控制。只取 1 个 batch完成一次前向反向然后直接退出。它能覆盖检测 Dataset 输出是否和模型输入对齐的大部分问题执行时间只有几秒。6. 数据管道的扩展与进阶方向当数据规模从 GB 涨到 TB6.1 用 MindRecord 封装大文件替代散装文本读取当语料规模上到几十 GB 甚至 TB 量级直接用TextFileDataset或GeneratorDataset读原始文本就有点撑不住了。原因是原始文本的 IO 开销大而且没有索引结构无法快速随机读取。昇思 MindSpore 提供了 MindRecord 格式相当于一种专门为训练设计的二进制存储格式。它会把数据按固定 chunk 写入.mindrecord文件并且自带索引数据读取和 shuffle 的效率比读文本高不少。writer ds.MindRecordWriter(corpus.mindrecord, files_num8) for text, label in preprocess_all_text(): writer.write({text: text, label: label}) writer.finish() # 训练时直接读 dataset ds.MindDataset(corpus.mindrecord, num_shardsrank_size, shard_idrank_id)用 MindRecord 之后shuffle的 buffer 可以适当缩小因为每个文件内部的数据已经足够乱序。这个形式在 PB 级数据处理团队里很常见数学上就是先粗粒度随机再细粒度随机的两级采样策略比纯全量 shuffle 在 IO 上便宜得多。6.2 数据并非越多越好先做样本量盘点再提效很多人在大模型数据预处理阶段容易陷入一个误区觉得数据越多越好所以把所有预处理都堆在 pipeline 里结果发现训练一个 epoch 要半小时其中大半时间都花在 map 算子的 tokenizer 转换上。我的经验是先做一次样本级耗时抽样吃透瓶颈在哪。用time模块或 profile 工具跑 1 万条样本的预处理耗时看 tokenizer 占多少、图像解码占多少、归一化占多少然后针对最耗时的部分做优化。如果 tokenizer 是大头最常见的手段是换成更轻量的 tokenizer或者使用缓存的分词结果让数据管道直接消费input_ids而不是每次训练时重新分词。如果无法绕开考虑扩大num_parallel_workers或者把 tokenizer 的结果写入 MindRecord训练时直接读。还有一个容易被忽略的点多次 padding 会浪费大量内存。如果你对固定长度 512 的样本已经做了 padding而数据集中实际长度小于 128 的样本占了 70%那相当于白白算了 4 倍的内存和注意力计算量。先统计长度分布再定max_length比盲目选择 512 要合理得多。对于大模型训练max_length256可能就是最经济的选择省下的显存能换更大的 batch size训练效率反而更高。6.3 pipeline 与模型训练交替迭代的小技巧最后分享一个更贴近开发节奏的技巧数据管道和模型结构其实是互相约束的。当你改动了模型结构比如把 attention 的 head 数改了input_ids的维度并不受影响但如果你改动了labels的 paddding 策略丢给 loss 函数的 label 就很容易发生 shape 不匹配。所以我一般会建立一套数据契约把每个字段的 shape 和 dtype 以常量形式写在一个公共模块里Dataset 管道和模型初始化都引用同一份定义。这样改一处的 shape另一处立刻能感知从根上避免了两边脱节的问题。# data_contract.py INPUT_ID_LEN 256 IMAGE_SIZE (224, 224) LABEL_ID_LEN 1这套思路虽然简单但能显著减少大模型微调时数据 shape 和模型输入对不上这类低级错误。我在多个项目里验证过最终沉淀下来也就这几条顺序要对缓存别乱加分布式种子要固定脏数据源头清洗。把这些做到位mindspore.dataset基本就能在正式训练开始前帮你挡住绝大多数数据问题。
网站建设高端定制企业官网