从 MySQL 到 JSONL:构建大模型数据加工流水线的工程实践

发布时间:2026/8/29 10:27:47
从 MySQL 到 JSONL:构建大模型数据加工流水线的工程实践 继 Seed、Flow 之后又被报道成立一个直接指向“数据”的 AI 一级部门这对做数据平台和基础架构的工程师来说可能比某一个大模型发布更值得关注。它背后释放的信号是大模型竞争已经不只是模型权重、推理算力和应用编排的竞争数据整理、数据治理和数据资产化正在变成一条独立的技术主线。对一线开发同学来说真正需要吃透的是“如何把关系数据库里的数据加工成大模型能读懂的数据”这条流水线。这篇文章不评述公司组织架构而是从工程视角拆解“数据”这件事为什么值得单独立起来以及在落地时通常要经历哪些步骤。文章会带大家完成一个最小可运行案例用 Python、Pandas 和 MySQL 造一张业务表再把它加工成可用于指令微调或评估的 JSONL 数据集。读完后可以拿着同样的思路回到自己项目里去搭一条最小数据流水线并且知道后面还要补哪些平台化能力。1. 为什么“数据”会成为独立的 AI 一级部门1.1 大模型的能力构成早已不只是模型结构很多人理解大模型时会把注意力全部放在 Transformer 架构、参数量、上下文长度这些模型侧能力上。但真正做过模型迭代的工程师都知道模型结构一旦定型决定效果上限的往往是数据预训练语料的质量、指令数据的覆盖度、偏好数据的一致性、评估数据的稳定性都会直接反应到模型表现上。这也是“数据”会从算法团队里拆出来的原因。过去一个部门里既有算法工程师又有数据工程师数据需求通常被当成边角料处理等模型训练前才临时清洗一批数据出问题后互相甩锅。当模型开始面向真实业务场景时标注、清洗、去重、脱敏、版本管理、血缘追踪这些工作会迅速膨胀如果继续混在算法团队里很容易变成瓶颈。独立的 AI 数据部门本质上是把数据当成一个独立的工程产品来建设。它的职责不是“帮算法团队跑几条 SQL”而是持续产出高质量、低风险、可追溯、可回滚的数据集。1.2 数据团队和算法团队的分工边界数据团队和算法团队如果边界不清最常见的现象是模型效果差的时候算法同学说是数据脏数据同学说是模型结构不行。两边各自都能找到理由但问题始终没有闭环。比较清晰的分工是算法团队负责模型结构、训练策略、推理优化、评估实验设计和模型发布。数据团队负责数据采集、清洗、格式化、质量检查、版本管理、敏感信息识别、数据血缘和发布。业务团队负责提供业务语义、标注规则、行业知识以及最终效果反馈。数据团队不是算法团队的附属而是一个独立交付方。它对外交付的不是一张表而是“一份可以用于建模的数据集”并且这份数据集要能回答三个问题数据来自哪里、经过哪些处理、当前版本和上一版本相比改了什么。1.3 从 Seed、Flow 到数据部门一条完整技术链路字节的 Seed 和 Flow 如果分别聚焦模型层和应用层那么数据部门实际上处在更底层的位置。数据部门不是在和模型团队并列而是在为整个技术栈提供共同的数据底座。这条链路可以简化成数据 - 模型 - 应用 - 用户反馈 - 数据。用户在使用应用时产生的反馈会重新变成数据回到池子里用于下一轮模型迭代。没有闭环的数据体系模型迭代就会变成“拍脑袋改数据”每次训练前都要重新梳理一遍成本极高。对中小团队来说不一定要单独立一个一级部门但必须有明确的数据 Owner。这个 Owner 要负责回答数据的口径是什么、质量怎么保证、线上发现数据问题后走什么流程修复。2. 从关系数据库到“大模型能读懂的数据”仍是一张复杂地图2.1 先区分训练数据、微调数据、评估数据和应用数据很多同学一听到“数据”第一反应就是把整张业务表导出然后丢给模型。这里最大的误区是没有区分不同用途的数据。大模型场景里至少需要区分五类数据数据类型主要用途典型格式最需要关注的指标预训练语料训练基座模型大规模纯文本去重率、毒性、信息密度指令微调数据SFT 阶段对齐指令JSONL包含 instruction、input、output指令覆盖度、答案一致性偏好数据RLHF 或 DPO 阶段使用包含多个回复及排序偏好一致性、标注者一致性评估数据模型效果回归固定评测集答案唯一稳定性、难度分布、防泄漏RAG 知识库数据应用阶段检索增强切块后的文本、向量和元数据切块质量、召回率、引用准确率同一张业务表可能被加工成不同类型的样本。比如用户反馈表可以提炼成“问题分类”的评估集也可以生成“客服回复”的微调数据还可以改写成“知识库问答”的 RAG 语料。使用方式不同加工逻辑完全不同。2.2 数据库里的数据为什么不能直接投喂给模型关系数据库里的数据是为事务查询和统计分析设计的不是为模型训练设计的。直接拿原始字段去喂模型会遇到几类问题第一字段语义不完整。比如一个订单表里有status3数据库里可能代表“已退款”但模型没有字段字典根本不知道 3 是什么意思。第二包含大量无关字段。用户 ID、创建时间、内部标记位对模型预测没有帮助反而会增加噪声。第三存在敏感信息。手机号、身份证、地址、账号往往直接存在业务表里直接进入训练集属于高风险行为。第四格式不稳定。同一个字段在不同时期可能写入不同格式比如日期有的是2024-01-01有的是20240101模型难以学到稳定规律。所以从关系数据库到模型可消费数据中间必须经过一层加工。这层加工不是简单复制而是要做语义补全、字段过滤、格式统一和敏感信息处理。2.3 数据加工流水线的标准步骤一条稳定的数据加工流水线通常包含六个环节抽取从 MySQL、Hive、Kafka 等数据源获取原始数据。清洗处理空值、去重、格式归一化、过滤异常值。格式化把业务字段拼装成模型需要的 instruction、input、output 结构。质检统计数据量、重复率、空值率、长度分布并抽查样本。版本归档按日期或版本号保存数据并记录生成代码和参数。消费供模型训练、模型评测或 RAG 应用使用。这六个环节缺一不可。跳过清洗后面的模型效果会不稳定跳过版本归档出问题后无法回滚跳过质检只能等模型训完才发现数据有问题。2.4 用一张表看懂数据管道涉及的组件如果只是本地实验一个 Python 脚本就能完成上面的流程。但到生产环境每个环节都需要对应组件支撑环节生产环境常用组件作用抽取Canal、Flink CDC、DataX从数据库或日志系统同步数据传输Kafka、Pulsar解耦上下游支撑实时数据转换Python、Pandas、Spark、dbt清洗、聚合、格式化存储MySQL、Hive、Iceberg、Delta Lake存储中间结果和最终数据集调度Airflow、DolphinScheduler定时触发、失败重跑、依赖编排目录DataHub、Atlas记录表结构、血缘、Owner监控Prometheus、Grafana任务状态、数据量、质量指标告警学习阶段可以用单机 Pandas 做完整过程但进入生产后靠人肉执行脚本是撑不住的。下面先实现一个最小案例让整条链路跑通。3. 最小可运行案例把 MySQL 业务表加工成指令数据集3.1 准备一张用户反馈表先建一张用户反馈表字段尽量贴近真实业务。这里用 MySQL 语法实际项目可以根据库表结构调整。CREATE TABLE service_feedback ( id BIGINT PRIMARY KEY AUTO_INCREMENT, user_name VARCHAR(64) COMMENT 用户昵称, category VARCHAR(32) COMMENT 反馈分类, content TEXT COMMENT 反馈内容, status TINYINT COMMENT 0-待处理 1-已处理, created_at DATETIME COMMENT 创建时间 ) ENGINEInnoDB DEFAULT CHARSETutf8mb4; INSERT INTO service_feedback (user_name, category, content, status, created_at) VALUES (张伟, 物流, 包裹已经三天没有更新了客服也不回复。, 1, 2025-01-12 10:00:00), (李娜, 售后, 申请退货之后一直没有人审核订单号是 A123456。, 0, 2025-01-12 11:30:00), (王强, 支付, 重复扣款了两次多扣的那笔什么时候退, 1, 2025-01-12 12:10:00), (NULL, 账号, 登录总是提示验证码错误换了好几个浏览器都一样。, 1, 2025-01-12 13:20:00);这张表里有用户昵称、分类、反馈内容和状态字段。用户昵称可能为空反馈内容里可能包含订单号这些在加工时都要考虑。3.2 用 Python Pandas 读取关系表本地环境需要安装依赖。推荐使用虚拟环境避免污染系统 Python。pip install pandas pymysql sqlalchemy读取数据库时推荐使用 SQLAlchemy 统一管理连接避免每次拼接数据库字符串导致编码问题。import pandas as pd from sqlalchemy import create_engine engine create_engine( mysqlpymysql://root:your_passwordlocalhost:3306/company?charsetutf8mb4 ) df pd.read_sql( SELECT id, category, content, status, created_at FROM service_feedback, engine ) print(df.head()) print(df.info())这里需要注意charsetutf8mb4一定要加。如果不加读取中文时经常会出现乱码尤其在生产库字符集不统一时。3.3 按业务字段拼接成训练样本要把业务表变成指令数据集关键是把category和content组合成一条完整的自然语言输入。下面是一个用于演示的拼接逻辑。df[instruction] 你是一名售后客服专家请根据用户反馈内容判断问题类型并给出建议。 df[input] ( 反馈分类 df[category].fillna(未分类) \n反馈内容 df[content].astype(str) ) # 演示用 output实际生产中应由标注结果或已有结论填充 df[output] 问题类型 df[category].fillna(未分类) 建议尽快联系用户核实。这里要特别强调不要把output简单用原始字段拼接。指令微调数据的output应当是高质量答案最好来自人工标注、线上已有客服回复或者经过审核的结论。如果只是用原始字段来回拼接模型学到的只是复制粘贴而不是理解和推理。3.4 清洗与去重数据处理中最容易出问题的是空值和重复。空值会直接导致模型训练时出现异常样本重复数据会让模型在局部样本上过拟合。# 去掉关键字段为空的行 df df.dropna(subset[content]) # 首尾空格清理 df[input] df[input].astype(str).str.strip() df[output] df[output].astype(str).str.strip() # 过滤过短的无效内容 df df[df[content].str.len() 5] # 生成去重标志 df[dedup_key] df[input] || df[output] import hashlib df[sig] df[dedup_key].apply( lambda x: hashlib.md5(x.encode(utf-8)).hexdigest() ) before len(df) df df.drop_duplicates(subsetsig, keepfirst) after len(df) print(f去重前样本数{before}去重后样本数{after})在生产环境去重不能只看精确字符串还要考虑语义重复。例如“包裹三天没更新”和“包裹已经三天没有物流记录”大概率是同一类问题但 MD5 去重无法识别。复杂场景需要引入 SimHash、MinHash 或向量相似度去重。3.5 输出成 JSONL 并做格式校验指令微调阶段最常用的格式是 JSONL每一行是一个完整 JSON 对象方便流式读取。df[[instruction, input, output]].to_json( sft_data.jsonl, orientrecords, linesTrue, force_asciiFalse )生成之后一定要做格式校验不能只看文件有没有生成。import json with open(sft_data.jsonl, r, encodingutf-8) as f: for line_no, line in enumerate(f, 1): try: obj json.loads(line) assert set(obj.keys()) {instruction, input, output} assert obj[input] assert obj[output] except Exception as e: print(f第 {line_no} 行格式异常{e})这段校验代码会捕获两种典型问题JSON 语法错误以及字段缺失或为空。把校验脚本放到流水线里就能在训练前发现问题。3.6 常见代码坑这个最小案例里已经藏着几个高频坑这里单独列出来。第一个坑是数据库连接缺少charsetutf8mb4导致中文读出来是乱码。检查方式很简单打印df[category]看是否正常然后确认建表语句和连接串的字符集。第二个坑是没有过滤空值。如果content是空字符串或者 NULL拼接出来的样本会非常短模型学到的是噪声。第三个坑是重复样本未消除。尤其是业务表里有大量内容相同的反馈时不做去重会让模型在训练时反复看到同一条样本导致记忆而不是泛化。第四个坑是 output 里带敏感信息。比如上面示例数据中的订单号A123456如果直接留在训练数据集里模型可能在推理时记住并泄露。4. 数据质量到底怎么度量4.1 质量不是“感觉干净”而是可量化数据质量不能靠“看起来还行”来评价需要落到一组可量化的指标上。常用的指标包括样本总数判断数据集规模是否满足训练需求。重复率重复样本占全量样本的比例。空值率关键字段为空的比例。字段长度分布过长或过短的样本占比。标签一致性同一条输入是否存在多个矛盾答案。敏感信息命中率手机号、身份证、邮箱等正则命中的比例。这些指标不需要一次性全部做但至少要在每次数据迭代时输出一份质量报告让模型团队可以对比版本差异。4.2 在管道里加质量检查可以把质量检查直接写进数据管道达不到阈值就终止后续训练流程。qa_report { total: len(df), empty_input_rate: round(float(df[input].isna().mean()), 4), duplicate_rate: round(1 - after / before, 4), avg_input_len: int(df[input].str.len().mean()), max_input_len: int(df[input].str.len().max()), } print(qa_report) # 如果重复率超过 10%说明清洗逻辑可能出现问题 if qa_report[duplicate_rate] 0.1: raise ValueError(数据重复率过高请检查抽取逻辑)这样做的意义是把质量判断从人的经验转为机器可执行的门禁。数据团队每次交付数据集时都能自动带出一份质量报告。4.3 用测试集衡量数据改动而不是只看模型 loss数据团队改数据不能只看训练 loss 是否下降。训练 loss 下降可能是因为模型记住了重复数据而不是真正学到了规律。更稳妥的做法是准备一份固定的评测集数据改动前后分别跑一遍评测集对比准确率、召回率或相关业务指标。评测集要保持稳定不要每次顺手改几道题。建议每个数据团队都维护“黄金评测集”数量不用太多几百条即可但要覆盖核心场景、边缘场景和对抗场景。数据变更后先跑黄金评测集再决定是否进入全量训练。5. 数据治理、安全和权限必须和开发同时进行5.1 敏感信息识别业务表里的敏感信息是训练数据集最常见的风险来源。手机号、身份证、邮箱、地址、订单号都可能被模型记住并在推理时输出。第一步是做好字段级规则识别。可以先从正则开始import re def mask_sensitive(text: str) - str: text re.sub(r1[3-9]\d{9}, [手机号], text) text re.sub(r\b\d{17}[\dXx]\b, [身份证], text) text re.sub(r[A-Za-z0-9._%-][A-Za-z0-9.-]\.[A-Za-z]{2,}, [邮箱], text) return text df[content] df[content].astype(str).apply(mask_sensitive)正则只能解决规则明确的敏感信息。真实场景里还有大量描述性敏感信息例如“我叫张伟住在上地三街 9 号”这类信息需要结合 NER 模型或人工抽检。比较好的策略是三层防线字段级规则、模型识别、人工抽检。5.2 数据版本和血缘数据加工不是一次性的模型会持续迭代业务表结构也会变化。如果不做版本管理下次想复现某个训练集就会非常困难。数据集命名建议采用数据集名称_日期_版本的格式例如sft_feedback_20250112_v1.jsonl。同时记录一份血缘信息数据来源哪张业务表、哪个库、哪个同步任务。处理逻辑清洗脚本版本、指令拼接代码版本。运行参数运行时间、调度任务 ID、责任人。变更说明本次相比上一版本改了什么。这些信息可以存在一个单独的表里和数据集文件放在一起。一旦训练效果异常可以快速回滚到上一个数据版本。5.3 数据资产目录当数据表数量变多后最怕的是没人知道某个表是干什么的。数据资产目录就是来解决这个问题的。一份基础目录至少要包含表名和用途说明。字段字典和口径解释。数据负责人。更新频率。下游消费方。没有目录之前新同学接手数据任务通常靠“问人”。有了目录后团队内部可以基于统一口径协作减少“同一个字段两种解读”的问题。6. 从脚本到平台规模化以后要补的工程短板6.1 缺少调度本地跑通了脚本并不等于生产可用。生产环境的数据任务需要调度系统来管理依赖关系、重跑逻辑和告警机制。比如数据加工任务依赖上游 ETL 任务先完成如果上游任务失败下游任务就不应该启动。调度系统可以用 Airflow、DolphinScheduler也可以用云厂商的托管调度服务。调度体系里至少要有三个能力定时触发、失败重试、超时告警。缺少任何一个线上任务都会变成“靠人盯”。6.2 缺少对象存储训练数据集通常很大不建议放在本地磁盘尤其是 K8s 集群里容器重建后文件可能直接丢失。生产环境建议使用对象存储按目录分区管理。s3://ai-data/sft/feedback/2025/01/12/sft_feedback_20250112_v1.jsonl对象存储的好处是容量可以扩展同时方便不同机器读取训练文件。学习阶段可以放到本机但一定要在代码里留出“把路径替换成对象存储”的抽象。6.3 缺少监控脚本定时跑起来之后最怕的是它“悄悄失败”。因此还要配置监控告警任务失败告警指定负责人收到消息。数据量跌幅告警如果某天样本数量比上周同期跌了 30%需要报警。质量指标告警重复率、空值率超过阈值时触发。敏感信息告警脱敏后仍有手机号命中时触发。监控是数据工程最重要的兜底。数据任务不像在线接口挂了不会有用户立刻投诉但会在训练环节积累成事故。6.4 字段变更业务表字段频繁变更是数据团队最大的隐性风险。某天发现训练数据全部为空很可能是因为源表字段名从content改成了feedback_content。生产环境建议创建数据契约也就是对输入表的字段结构做校验required_columns [id, category, content, status, created_at] missing [c for c in required_columns if c not in df.columns] if missing: raise ValueError(f源表缺少字段{missing})字段契约可以在每次任务启动时校验一次提前发现问题而不是等训练跑到一半才报错。7. 常见问题与排查清单7.1 现象、原因、检查方法、处理建议下面这张表收集的是数据加工场景里出现频率最高的问题可以直接作为排查手册使用。问题现象常见原因检查方式处理建议训练数据里中文乱码数据库连接缺少 charset或源表字符集不一致打印样本字段检查连接串和表字符集统一使用 utf8mb4并在连接串显式指定JSONL 解析失败文本中包含换行符或未转义引号手工拼接 JSON用 json.load 逐行解析定位异常行号统一用 json.dumps 序列化不手工拼字符串模型回答总重复历史内容数据重复率过高训练集里大量相似样本统计精确重复率和语义相似度引入 MinHash 或向量去重训练后效果反而下降新数据过滤规则过于激进删掉关键样本对比过滤前后的数据分布和评测集结果先用小批量数据做消融实验数据集里有用户手机号脱敏逻辑未覆盖全部文本字段用正则扫全量文本抽查输出文件增加字段级规则和人工抽检某天数据集样本数为 0源表字段变更或上游同步任务未执行检查字段契约并查看上游任务日志增加字段校验和调度依赖7.2 排错链路遇到数据处理问题不要一上来就怀疑模型也不要一上来就重跑全量任务。按下面顺序排查效率更高先确认输入数据是否正常查询源表行数、字段内容、时间范围。确认脚本读取阶段是否报错打印 DataFrame 的 shape 和前几行。确认清洗逻辑是否把合法数据过滤掉了对比清洗前后的数据量。确认输出文件是否能被正常解析逐行校验 JSON。确认质量报告中的指标是否合理重复率、空值率是否异常。确认调度任务上下游是否都成功查看任务依赖、重试日志。数据链路的问题90% 都可以在步骤 1 到 3 里定位。如果这三步都没发现问题再考虑上游同步任务或数据库连接不稳定。8. 最佳实践一个 AI 数据团队的核心工作方式8.1 数据不是一次性产物需要版本化很多团队做数据是这样模型训练前算法同学扔给数据同学一个需求数据同学连夜写脚本导出 CSV训完就丢。下次需要时又从头开始。正确做法是让每个数据集都像代码一样有版本。版本里至少包含数据集名称和版本号。生成时间。来源表和时间范围。清洗脚本和参数。质量报告。变更说明。有了版本后模型复现和问题回溯才成为可能。8.2 模型团队、数据团队、业务团队怎么配合数据团队不能闭门造车也不能等到模型团队提需求才动手。比较好的协作方式是建立“数据需求单”流程业务团队提出业务场景和口径。算法团队提出数据格式和样例要求。数据团队负责实现并交付质量报告。算法团队在评测集上验证并反馈结果。数据团队根据反馈迭代数据版本。在这个流程里数据团队是独立交付方而不是“临时帮工”。8.3 给工程师的落地方案如果你所在团队还没有成型的数据基础设施建议按下面的节奏推进先跑通一条最小数据流水线只覆盖一个业务场景。在流水线里加入质量检查输出质量报告。把数据集保存成带版本和血缘信息的文件。接入调度任务和失败告警。完善敏感信息识别和权限管理。再逐步接入自动化标注、主动学习、向量去重等高级能力。前四步可以在一到两周内完成但能显著减少“数据事故发生在训练完成之后”的概率。第 5、6 步才是长期优化方向。回到开头提到的公司新闻字节把数据独立成 AI 一级部门本质上是在组织层面承认“数据是模型能力的核心变量”。对普通工程师来说这件事同样有参考价值无论公司有没有独立数据部门你参与的项目都应该把数据加工、数据质量、数据版本和数据安全当成一等公民来对待。先从最小案例开始把整条数据链路跑通再逐步完善平台化能力这是最稳妥的路径。