
1. 大模型训练里数据预处理为什么排在模型结构前面做昇思MindSpore大模型训练很多朋友把精力全放在网络结构、loss 设计上结果一跑起来第一个卡住的反而是数据加载和预处理。尤其是从 PyTorch 转过来的人写完 Dataset 后顺手在__getitem__里做分词、padding、mask到了 MindSpore 里发现GeneratorDataset的行为不一样mindspore.dataset的 map、batch、shuffle、repeat 顺序也讲究一上来就懵。这篇文章我打算把基于mindspore.dataset的数据变换与预处理讲透重点覆盖大模型微调和预训练里最常见的文本处理场景包括 JSONL 原始数据加载、分词、token 化、padding、attention_mask、labels 构造、变长 batch以及多卡并行和性能调优。适合正在做微调、RAG、指令跟随模型或者从头训练小尺寸 LLM 的朋友参考也适合想从 PyTorch 迁移到昇思的算法工程师快速上手。先说一个核心观念数据预处理不是把文件读进来交给模型就完事它是在训练开始之前把你的文本、图像、标签统一成模型需要的张量形状和数值语义。大模型尤其在意两件事一个是 sequence length 对齐因为 transformer 的 attention 机制要求 batch 内 token 数量一致另一个是 padding 位置不能参与 loss 计算否则模型会学到“输出 pad token 也没关系”最后生成质量崩掉。这两件事必须在mindspore.dataset管道里解决而不是散落在训练循环里临时处理。1.1 数据管道不是“搬数据”而是在训练开始前完成特征约定mindspore.dataset给我最大的感受是它把“读数据-变换-组 batch-训练迭代”做成了流水线。你定义好GeneratorDataset、MindDataset、TextFileDataset后后续的map、batch、shuffle、repeat其实都是惰性操作真正执行是在create_dict_iterator或者model.train被调用时才触发。这种惰性机制有个好处就是你可以在训练前先用一个小循环检查 pipeline 输出对不对而不用把整个数据集全部加载到内存。坏处是一旦你不理解 map 的执行时机很容易写出“看起来对跑起来错”的代码。比如我见过有人在GeneratorDataset里直接读文件、分词、padding、返回固定形状数组这当然也能跑但会丢掉mindspore.dataset自带的 shuffle、并行、动态分桶能力数据加载速度会差很多。更合理的做法是GeneratorDataset只负责按行读取原始样本tokenizer 之类的重活交给map算子shuffle放在 map 之前或之后根据数据量决定batch和 padding 放到最后repeat包在最外面。这样每一步职责单一出问题也好排查。1.2 最容易犯的错把预处理写进模型前向还有一类朋友会在模型construct里做 token 化比如把字符数组传进网络再在算子里调用分词。这个在数据量小的时候看不出问题数据量一大代价非常恐怖。大模型训练几乎都是 batch 级矩阵运算GPU/NPU 最怕的就是在训练过程里插入不可并行的 Python 循环。分词器是 Python 对象逐个 sample 跑编码速度可能比前向本身还慢。昇思官方建议把数据变换尽量下沉到 dataset 层原因就在这里mindspore.dataset的map可以用多进程并行处理还能和计算设备重叠执行。我可以负责任地讲一次 SFT 微调里如果 tokenizer 放在数据管道外10 万条样本可能要额外多花 3 到 5 个小时放到 pipeline 里并用num_parallel_workers8这段耗时能被压到几十分钟。所以凡是能在数据准备阶段算好的东西绝对不要留到训练阶段再算。2. mindspore.dataset 基础组件与流水线顺序2.1 数据加载器选型mindspore.dataset提供了一堆 Dataset 类大模型场景里常用的主要是这几个加载器适用场景我的使用建议GeneratorDataset自定义 Python 生成器、从 JSONL/CSV 读入、无法用现成 Dataset 表达的数据最灵活微调场景首选MindDataset读取 MindRecord 文件预处理后落盘多 epoch 重复训练时首选TextFileDataset每行一条纯文本的语料做词表、预训练语料清洗时可以用TFRecordDataset已经从 TF Record 转换好的数据迁移旧数据集时用ImageFolderDataset按文件夹组织的图片分类数据多模态或视觉任务里常见大模型微调最常见的输入是 JSONL每行是一个样本包含 prompt、response甚至可能还有 chosen/rejected 用于偏好对齐。这种结构没有现成 Dataset用GeneratorDataset是最顺的。用GeneratorDataset时有几个参数必须关注column_names决定返回字段名num_shards和shard_id多卡并行时一定要设置shuffle不设置的话很多自定义生成器默认是不打乱的。我踩过多次坑就是忘记传shuffleTrue导致多卡每张卡拿到的数据分区相同训练曲线看似正常实际评估一塌糊涂。2.2 map、shuffle、batch、repeat 的先后顺序流水线顺序不是随便排的。我的习惯是原始数据先 map 做 tokenizer 等单样本变换shuffle 打乱顺序buffer_size 尽量大于数据集中重复模式长度再 map 做 padding 和 mask、labels 这类依赖 token 结果的操作batch 组装成 batchrepeat 在最外层控制 epoch。为什么 tokenizer 要先于 shuffle因为大模型训练要求样本随机性足够强如果你先 batch 再 shuffle那 shuffle 的只是 batch 顺序batch 内部样本顺序仍然是固定的会让梯度方差偏离预期。为什么 padding 要放在 batch 前因为如果不先 pad 到固定长度batch算子拼接出来的张量形状不一致MindSpore 会直接报错。repeat包在最外层还有另一个原因model.train(epochs, dataset)在昇腾环境开了dataset_sink_modeTrue时MindSpore 自己会迭代数据。如果你在 pipeline 外层额外 repeat有时会出现 epoch 次数不对、数据提前耗尽等问题。我的建议是常规训练里直接用data.batch(...).repeat(epochs)交给model.train消费或者在model.train外层控制 epoch不要在 pipeline 里和 train 参数里同时设置 epoch。2.3 transforms 也要分 C 算子和 Python 算子mindspore.dataset的 transform 分了两类一类是内置算子比如Resize、PadEnd、Lookup、BertTokenizer底层是 C 实现速度快另一类是你自己写的 Python 函数通过map(operationsfunc)传入底层按 Python 调用执行。这个区分直接决定性能。内置 C 算子单线程处理能力也很强配合num_parallel_workers可以跑满 CPUPython 函数受 GIL 限制多线程效果有限要开python_multiprocessingTrue才能真正并行。所以设计 pipeline 时能用内置text.Truncate、text.PadEnd、text.Lookup、ds.vision.Resize这类算子就不要自己写 Python 循环。只有像 HuggingFace tokenizers 这种无法绕过的第三方库再考虑用 Python 函数包一层并且一定要开 multiprocessing。3. 通用文本数据变换与预处理完整实战3.1 从 JSONL 到 GeneratorDataset下面这段代码是大模型 SFT 最常用的起点。假设你的 JSONL 长这样{prompt: 解释一下反向传播, response: 反向传播是一种计算梯度的算法...} {prompt: 写一首五言绝句, response: 明月出天山苍茫云海间。}读取和构造 Dataset 的代码import json import numpy as np import mindspore.dataset as ds from tokenizers import Tokenizer tokenizer Tokenizer.from_file(my_llm_tokenizer.json) def load_samples(path): with open(path, r, encodingutf-8) as f: for line in f: line line.strip() if not line: continue obj json.loads(line) yield obj[prompt], obj[response] data ds.GeneratorDataset( load_samples(train.jsonl), column_names[prompt, response], shuffleTrue, num_shardsrank_size, shard_idrank_id )这里有几个细节值得展开。GeneratorDataset的 source 参数传的是一个生成器函数对象不是已经构造好的列表。如果直接传load_samples(train.jsonl)MindSpore 会在内部迭代时反复调用这个生成器很容易出现文件句柄被二次消费的问题。正确写法是传load_samples这个可调用对象让 dataset 自己管理迭代生命周期。column_names必须和yield的顺序一一对应。yield 两个字段column_names 就写两个名字。多卡训练时num_shards和shard_id通常从环境变量取rank_size int(os.getenv(RANK_SIZE, 1)) rank_id int(os.getenv(RANK_ID, 0))不设置这两个参数单卡没问题多卡就会数据重复。3.2 分词、词表映射与长度截断拿到原始文本后下一步是 token 化。工程上我推荐直接用 HuggingFace tokenizers 或 SentencePiece因为现在绝大多数开源大模型都已经训练好了 BPE/Unigram 词表你用自己的BasicTokenizer重新分词很可能和模型预训练时不一致导致 embedding 分布错位。用tokenizers库做处理包成一个 Python 函数prompt_template sHuman: {}\nAssistant: {} def encode_sample(prompt, response): text prompt_template.format(prompt, response) ids tokenizer.encode(text).ids return np.array(ids, dtypenp.int32) data data.map( operationsencode_sample, input_columns[prompt, response], output_columns[input_ids], num_parallel_workers8, python_multiprocessingTrue )这里我特意用了两个输入列演示map的多列输入用法。函数签名里的参数顺序必须和input_columns一致。返回结果用 numpy 数组MindSpore 会统一转成 Tensor。如果不想用外部库而是基于昇思内置的文本 transform 自己搭词表可以这么做from mindspore.dataset import text # 先用一个浅管道拿到切词后的文本 word_ds ds.TextFileDataset(corpus.txt, shuffleFalse) word_ds word_ds.map(operations[text.BasicTokenizer()], input_columns[text]) # 从 dataset 构建词表 vocab text.Vocab.from_dataset( word_ds, columns[text], special_tokens[[PAD], [UNK]] ) # 再对训练集做词表映射 data data.map( operations[text.Lookup(vocab, unknown_token[UNK])], input_columns[input_ids], output_columns[input_ids] )注意text.Lookup期望输入是 token 序列如果你的数据已经是整数 id就不需要再做 Lookup。这个流程更适合从小规模语料训练自己的 tokenizer 和语言模型。分词之后长度截断必须放在 padding 之前。我习惯在同一个 map 里先text.Truncate再text.PadEndseq_len 512 def truncate_and_pad(ids): return text.Truncate(seq_len)(ids), text.PadEnd(pad_token_id, seq_len)(ids)不过更常见的做法是直接用内置算子链式 mapdata data.map(operations[text.Truncate(seq_len)], input_columns[input_ids]) data data.map(operations[text.PadEnd(pad_token_id, seq_len)], input_columns[input_ids])如果你用的是外部 tokenizer也可以直接在encode_sample里设置truncationTrue然后调用text.PadEnd做定长。截断长度要和你模型训练时的seq_length对齐别模型配了 2048数据管道还在用 512。3.3 padding、attention_mask 与 labels 的构造对齐后的 input_ids 还不能直接进模型。第一个要补的是 attention_mask让模型知道哪些 token 是真实内容哪些是 padding。最稳妥的做法是在 padding 之后用一个 Python map 函数生成 mask 和 labelsdef make_mask_and_labels(ids): ids np.asarray(ids) mask (ids ! pad_token_id).astype(np.int32) labels ids.copy() # 把 padding 位置设为 -100PyTorch 系 loss 会忽略MindSpore 的 CrossEntropyLoss 同样支持 ignore_index labels[mask 0] -100 return mask, labels data data.map( operationsmake_mask_and_labels, input_columns[input_ids], output_columns[attention_mask, labels], num_parallel_workers8, python_multiprocessingTrue )这里有一个容易忽略的细节如果 tokenizer 带了 special token比如s、/s、[CLS]、[SEP]attention_mask 要不要把这些也标成 1要看你的模型设计。绝大多数 LLM 的 causal attention 里special token 也是真实输入应该参与 attention只是 loss 计算时一般要把 padding 位置和特殊起始标记排除掉。labels 的构造同样要看模型实现。有的模型 forward 里已经做了 input shift也就是用input_ids[:, :-1]预测input_ids[:, 1:]那你的 labels 保持和 input_ids 对齐就行。有的模型没有做 shift需要你在数据侧手动把 labels 左移一位最后一个 token 补 paddef make_clm_labels(ids): ids np.asarray(ids) labels np.copy(ids) labels np.concatenate([labels[1:], [pad_token_id]]) labels[labels pad_token_id] -100 return labels对 SFT 指令微调来说更精细的做法是只对 response 部分的 token 计算 lossprompt 部分全部置成 -100。这样模型不会去背题只学习如何组织回答。具体做法是记录 response 在拼接后的起始位置在这个位置之前全部 mask 掉。这个信息你在构造 prompt 拼接字符串时是知道的建议直接存成一个新字段。3.4 动态长度分桶 batchbucket_batch_by_length很多微调数据的长度方差很大。有的样本只有 50 个 token有的样本 800 个 token。如果统一 pad 到最长长度短样本白白浪费算力如果统一 pad 到 512长样本又会被强行截断。昇思提供了bucket_batch_by_length它会根据样本长度把样本分到不同桶里再对不同桶用不同 batch size。比如短样本可以放 32 条一批长样本只放 8 条一批最后每个桶内部 pad 到本桶最大长度避免全部 pad 到全局最大。使用方式大致如下data ds.GeneratorDataset(load_samples(...), column_names[prompt, response]) data data.map(operationsencode_sample, ...) data data.bucket_batch_by_length( column_names[input_ids], bucket_boundaries[128, 256, 512], bucket_batch_sizes[32, 16, 8], pad_to_bucket_boundaryTrue, pad_info{ input_ids: (None, pad_token_id), attention_mask: (None, 0), labels: (None, -100) }, drop_remainderTrue )注意pad_info里 shape 可以写None表示自动按当前桶的最大长度填充。如果你的 MindSpore 版本不支持None就换成固定[512]效果是每个桶内部还是统一 pad 到 512只是 batch size 按长度做了分档。bucket_batch_by_length带来的收益非常直观。我有一份 SFT 数据长度中位数只有 256因为很多短问答如果硬 pad 到 1024浪费的显存接近 60%。换成桶 batch 后同样的 batch_size 设定训练吞吐提升了接近一倍。代价是每个 bucket 之间的样本顺序被打乱如果你的 loss 曲线需要严格可复现别忘了设置随机种子。4. SFT/指令微调场景下的数据变换进阶4.1 SFT 对话模板与多轮拼接指令微调的核心难点不在模型而在把多轮对话拼成模型期望的格式。常见的格式是sHuman: ...\nAssistant: .../s如果是多轮还要把历史轮次按顺序拼接。这里有个坑不能简单把每条对话 logger 过一遍 tokenizer 再拼接 token id因为很多 BPE 分词器对换行拼接位置敏感分开 encode 再 concat拼接处的 token 可能和整体 encode 不一样。正确的做法是先把完整对话模板组装成字符串再一次性 encode。比如def format_multi_turn(messages): # messages: [{role: user, content: ...}, {role: assistant, content: ...}] parts [] for msg in messages: if msg[role] user: parts.append(sHuman: msg[content]) elif msg[role] assistant: parts.append(\nAssistant: msg[content]) parts.append(/s) return \n.join(parts)在 JSONL 里可以用messages数组表示多轮也可以用history字段。无论哪种都要保证模板里的特殊 token 和你在训练时用的一致。这个拼接逻辑最好写单元测试因为文本预处理 bug 是隐性的模型可能还能训练但生成质量永远不对。4.2 标签掩码与 ignore_indexSFT 里最常见的 loss 错误是模型把 prompt 部分也纳入预测目标导致模型学会了复述问题而不是回答问题。很多人在数据管道里只做了 padding mask没有做 prompt mask最后 loss 居高不下。我的做法是额外存一个response_start字段或者直接用布尔向量标记 response 区域。举例说明def make_sft_labels(input_ids, response_start): input_ids np.asarray(input_ids) labels input_ids.copy() labels[:response_start] -100 return labelsresponse_start可以从模板拼接后字符串里计算出来也可以让 tokenizer 返回offsets来确定。最简单的方式是在format_multi_turn里返回(text, response_start)两个值。这个字段在 token 化之后作为第二个输入传给 labels 构造函数。MindSpore 的nn.CrossEntropyLoss不一定支持二维 logits 直接对应 sequence 的 ignore_index很多 LLM 实现会在 loss 计算前把 logits 和 labels 展开。无论用哪种实现数据侧的 labels 必须保证 padding 和 prompt 位置都是 -100这是底线。4.3 并行训练与 MindRecord 持久化多卡并行时我建议把预处理结果写进 MindRecord而不是每次都从 JSONL 重新分词。MindRecord 是昇思的原生二进制格式加载速度远快于文本文件而且天然支持分布式分片。大概流程是先用 pipeline 处理成input_ids/attention_mask/labels再落盘。具体代码如下from mindspore.mindrecord import FileWriter writer FileWriter(train.mindrecord, shard_num8) writer.add_schema({input_ids: {type: int32, shape: [-1]}, attention_mask: {type: int32, shape: [-1]}, labels: {type: int32, shape: [-1]}}, llm_sft) for item in data.create_dict_iterator(num_epochs1): writer.write_row({input_ids: item[input_ids].tolist(), attention_mask: item[attention_mask].tolist(), labels: item[labels].tolist()}) writer.commit()之后训练时用ds.MindDataset(train.mindrecord, num_shardsrank_size, shard_idrank_id)加载即可。这个做法最大的好处是分词这种 CPU 密集型操作只需要跑一次后续多轮实验只是在读 MindRecord省下的时间非常可观。5. 多模态模型里图像预处理怎么接进同一管道5.1 图像读取与增强大模型不止文本多模态模型现在也很常见。mindspore.dataset同样支持图像变换最基础的是Decode - Resize - Normalizefrom mindspore.dataset import vision data ds.GeneratorDataset(load_multimodal_samples(), column_names[image_path, text]) data data.map(operations[vision.Decode()], input_columns[image_path]) data data.map(operations[ vision.Resize((224, 224)), vision.Rescale(1.0 / 255.0, 0.0), vision.Normalize(mean[0.485, 0.456, 0.406], std[0.229, 0.224, 0.225], is_hwcTrue) ], input_columns[image_path])注意Decode之前image_path是字符串Decode后变成 uint8 HWC 数组。昇思的图像 normalize 默认是 CHW 还是 HWC 取决于版本老版本很多算子要求 HWC新版本有变化。如果你发现 shape 不对记得先加一个vision.HWC2CHW()做通道转换。5.2 文本-图片字段对齐多模态数据往往来自不同来源文本和图片单独处理到 batch 时字段拼接必须保持一致。最佳做法是让GeneratorDatasetyield 的每个样本本身就是(image_array, text_ids)对齐的文本 token 化在map中完成图像增强在另一个map中完成最后 batch 时两个字段同批次存在。如果图像路径解析失败比如某些样本文件缺失不要让它炸掉整个训练。在生成器里用白名单过滤掉空样本或者在map函数里对异常情况返回None并用filter算子剔除。mindspore.dataset的filter可以在 pipeline 里按条件筛数据比训练循环里判断要干净得多。6. 性能调优与实测问题排查6.1 常见性能瓶颈我实测下来mindspore.dataset处理大模型数据的性能瓶颈通常集中在三个地方第一是 Python 生成器本身。GeneratorDataset每次迭代都要把控制权交回 Python如果生成器里做了文件 IO 和 JSON 解析整条流水线都会卡在这里。建议生成器只做 IO 和解析把更重的分词交给map并行执行。第二是python_multiprocessingFalse时Python transform 串行执行。像 tokenizers 这类库单线程 encode 1 万条样本可能要几分钟开python_multiprocessingTrue后能按 CPU 核数翻倍。不过多进程有数据序列化开销样本体量特别小的时候不一定更快需要实测。第三是 shuffle buffer 太大或太小。buffer 太大内存占用高太小随机性不足。我习惯设成max(10000, 单个epoch样本数//1000)然后观察训练收敛速度和 loss 曲线波动性。还可以全局配置import mindspore.dataset as ds ds.config.set_prefetch_size(16) ds.config.set_num_parallel_workers(8) ds.config.set_enable_autotune(True) ds.config.set_seed(42)set_enable_autotune(True)是昇思的自动调优开关它会让系统根据实际运行情况调整并行度和 prefetch 大小。在新版本里比较稳老版本偶尔会出现 batch 边界问题如果遇到异常输出先关掉 autotune 再排查。6.2 问题速查表现象原因处理方式多卡训练总 loss 异常每卡数据重复创建 Dataset 时没传num_shards/shard_id补上分布式参数重新加载数据token 化之后长度不一致batch 报错没有截断或 padding在 batch 前加Truncate和PadEndloss 一直不降生成文本里出现 pad tokenlabels 中 padding 位置没置 -100检查 labels 构造函数padding 位置要 ignore数据 pipeline 跑得慢CPU 利用率低Python transform 未开多进程map加python_multiprocessingTrue适当加大num_parallel_workers数据顺序与预期不符复现困难shuffle 随机种子未固定全局set_seed并在每个数据集采样器设置 seed长文本被截断后回答被切断Truncate长度太短调整 seq_len 或启用 chunk 拆分逻辑使用bucket_batch_by_length后训练不稳定bucket 间样本分布不均增大 bucket 数量让长度分档更细必要时固定随机种子另外有一个我踩过多次的坑map的input_columns如果同时包含多个列Python 函数必须写成多参数形式而不是接收一个 dict。一旦写成def func(data):这种接收字典的形式昇思会按列顺序把每个字段当独立参数传入函数签名对不上就会报错。如果你确实想传递整个样本可以把多个字段合并成一个列再在 map 里操作。还有个容易忽略的细节create_dict_iterator(num_epochs1)是调试 pipeline 的神器但千万别在model.train之前反复调用它来“预热”。每次调用都会完整遍历一遍数据集不仅浪费时间还会改变某些带有随机性的 shuffle 状态。调试用data.take(5)就够了。最后分享一个我对mindspore.dataset的整体判断它不是一个需要死记 API 的东西核心是理解“数据流”。写出一个高效、稳定、可复现的预处理管道比背下一百个 transform 名重要得多。我的习惯是所有新数据集先写一个 10 分钟的 smoke test用take(2)打印input_ids/attention_mask/labels的形状和一个样本的 decoded 文本确认无问题后再铺到全量训练。这个习惯帮我避免了至少十次“训练到凌晨两点才发现数据错误”的尴尬。