Text-to-SQL 微调项目里最容易被低估的环节一定是数据管线。模型结构可以抄现成的Loss 可以调但训练集里的 SQL 质量直接决定模型能不能在真实业务库上写出能跑的查询。我最近搭了一套基于 DataFlow 的 Text-to-SQL 数据处理 Pipeline把线上审计日志、手工沉淀 SQL、开源数据集混合清洗成可直接送进 SFT 阶段的训练样本整个链路分成采集、解析、清洗、验证、转换五段每一段都有独立的 quality check。这篇文章把这条管线的设计思路、核心节点实现和踩过的坑完整复盘一遍适合正在做 SQL 大模型微调、或者想给自己的 RAG 系统补 Text-to-SQL 能力的朋友参考。先说清楚一个前提为什么不能直接拿网上的 Text-to-SQL 数据集来训练Spider、BIRD 这些开源基准质量确实不错但 schema 规模和业务习惯跟真实场景差异很大。真实线上 SQL 往往带着多层嵌套、奇怪的缩写命名、历史遗留字段甚至还有WHERE 11这种动态拼接痕迹。模型在开源数据上学得再漂亮落到自己业务库上往往生成一些“语法正确但逻辑跑不通”的 SQL。所以别绕开数据准备这一步它值得你投入比模型调参更多的时间。1. 先弄清楚Text-to-SQL 训练数据到底要什么1.1 为什么通用问答语料不够用大模型训练里通用语料负责“语言能力”Text-to-SQL 数据负责“结构化查询能力”这两件事没有天然的迁移关系。你让模型读一万篇技术博客它能学会 SQL 语法但学不会“把‘上个月每个区域的销售目标达成率’翻译成一段带多层 JOIN 的聚合查询”。Text-to-SQL 训练样本的本质是三元组对齐数据库 Schema、自然语言问题、SQL 语句三者必须严格一致。Schema 相当于“规则说明书”自然语言是“用户意图”SQL 是“执行方案”。缺了 Schema模型不知道表名和字段名缺了自然语言模型没法学会语义映射缺了 SQL模型无从模仿查询结构。这也是为什么通用问答语料无法替代专门数据——它们根本没包含 Schema 上下文。1.2 高质量 SQL 训练数据的五个特征我在清洗管线里把“高质量”拆成了五个可量化指标一是可执行性每条 SQL 必须能在对应 Schema 上真实跑通而不是语法上说得过去。语法正确和可执行是两码事表不存在、字段类型不匹配、隐式转换报错都是模型常犯的错误训练数据里如果混入这种坏样本模型会被带偏。二是语义一致性自然语言和 SQL 的逻辑必须一一对应。用户问“销售额”SQL 里不能只查了订单量“最近7天”这样的时间限定必须体现到 WHERE 条件里。三是结构多样性不同 SQL 的算子分布要合理。单表查询、多表 JOIN、子查询、窗口函数、CTE 都要有一定比例否则模型只会抄最常见的模板。四是Schema 覆盖度训练集涉及的表要足够多最好是全量业务表的 80% 以上不然模型遇到没见过的表就直接放弃。五是真实分布匹配真实业务 SQL 往往带着冗余条件、隐式类型转换、奇怪的别名这些“脏乱差”特征反而能让模型在推理时更健壮。但要过滤掉纯噪音比如调试语句、超长死循环查询、明显注入特征。2. 数据管线整体设计从原始SQL到训练样本2.1 数据源盘点这些 SQL 从哪来一个合格的数据管线输入源绝对不能只有一种。我目前接了四个来源数据库审计日志MySQL 开启 general log 或者使用云厂商的 SQL 审计能力能拿到线上真实查询。这是最有价值的数据源但噪音也最大DBA 的运维查询、ORM 框架自动生成的查询、监控探针查询全都混在里面。数据仓库查询历史Hive、Presto、Doris 等引擎的审计日志能拿到大量复杂分析型 SQL对提升模型处理大数据量查询的能力很有帮助。Git 仓库中的 .sql 文件数据工程师沉淀的脚本、报表系统的取数逻辑这类 SQL 结构完整、注释清楚适合作为“干净样本”的种子。人工修正样本让数据分析师把日常业务问题写成 SQL 对产出少量“标准答案”。量不用大但一定要有它们是验证整个管线质量的黄金标准。这些数据源格式完全不同有的是文本日志、有的是结构化文件、有的需要连库读取。DataFlow 里的 Resource 对象就是干这个的每种数据源声明成一个资源Pipeline 只关心资源的读写协议不需要关心数据从哪来。2.2 五个核心处理阶段的拆解每个阶段解决一个独立的问题阶段之间不互相依赖单点失败不会导致全链路崩溃。采集阶段负责把不同来源的数据统一成一种中间格式我用的是一份 JSON 记录包含原始 SQL 文本、来源标签、采集时间、所属库名。统一格式的好处是后续所有处理节点只需要面对一种输入。注意保留来源标签后面采样的时候有大用。解析阶段的核心任务是把 SQL 文本变成 AST。用 sqlglot 或 JSqlParser 这类语法解析器做词法和语法分析提取出 SQL 类型、涉及的库表、字段、JOIN 关系、WHERE 条件、使用的算子。解析阶段同时做方言归一化线上 MySQL 写的 SQL 和数据仓库的 Hive SQL 语法有差异统一转成标准 SQL 风格后续模型学习时不容易被方言细节干扰。清洗阶段做过滤和去重。过滤规则包括只保留 SELECT 查询、剔除 DDL/DML、剔除带注入特征的语句、剔除单条超过 2000 字符的超长 SQL、剔除嵌套超过 5 层的“俄罗斯套娃”查询。去重不是做文本比对而是做 AST 结构签名比对这个后面细说。验证阶段解决“这条 SQL 到底能不能跑”的问题。规范的做法是连一个影子库执行 EXPLAIN 或者 LIMIT 1 探测语法错误直接丢弃表不存在直接丢弃执行时间超过阈值直接丢弃。这个阶段是整个管线里最消耗资源的一环但绝对不能省略。转换阶段把通过验证的 SQL 加工成训练样本生成 Schema DDL、生成自然语言问题、打难度标签、按模板格式化输出为 JSONL。这一步的输出才是真正送到训练任务里的数据。2.3 用 DataFlow 编排后的整体链路DataFlow 的核心抽象是 Resource、DataFlow、Engine 三个概念。Resource 定义数据从哪里来、到哪里去DataFlow 定义数据的处理流向和每个节点的处理函数Engine 负责实际调度执行本地调试时用线程池引擎全量跑批时切到 Flink 引擎节点代码不用改。整个链路在 DataFlow 里长这样先声明四个 Resource 作为输入源再接一个sql_collector节点做统一格式转换然后依次是sql_parser、sql_filter、sql_dedup、sql_validator、sample_generator最后用一个dataset_writer节点把结果写入训练数据存储。每个节点都是无状态的处理函数输入输出通过 DataFlow 自动映射好处是单个节点可以独立重跑、独立扩并发某一个环节出问题不需要把整个链路重来一遍。这套编排方式还带来一个额外好处离线批处理和增量流式处理可以共用同一套节点逻辑。比如想要每天凌晨两点全量跑一次历史数据同时白天每十分钟增量处理新日志只需要换一个 Engine、改一下资源监听方式节点代码零改动。3. 核心实现SQL 清洗与增强的几个关键手段3.1 用 sqlglot 做语法解析和 AST 标准化sqlglot 是我在解析阶段的主力工具它比 Antlr 那套方案轻量得多而且对多种方言有不错的兼容性。核心用法很简单import sqlglot raw_sql select a.id, b.name from users a join orders b on a.id b.user_id where a.status 1 ast sqlglot.parse_one(raw_sql, readmysql) # 提取 SQL 类型 if isinstance(ast, sqlglot.exp.Select): print(是 SELECT 查询) # 提取涉及的表 tables {table.name for table in ast.find_all(sqlglot.exp.Table)} print(tables) # 提取 JOIN 条件 joins ast.find_all(sqlglot.exp.Join) print(len(joins)) # 标准化输出 normalized_sql ast.sql(dialectmysql, prettyTrue) print(normalized_sql)AST 标准化是这里最有价值的一环。同一个查询可以写成SELECT * FROM users WHERE id1也可以写成selectusers.* from USERS where id 1文本完全不同但 AST 结构几乎一致。我们在 AST 层做归一化所有标识符转小写、去掉反引号、字符串字面量替换成占位符、去掉多余括号和冗余别名最终得到一份“标准形态”。这份标准形态既是去重的基础也是后续难度分级和模板聚类的输入。实操中有个经验不要试图用一个方言参数解析所有 SQL。线上日志里的 SQL 经常混着奇怪写法比如 MySQL 的ON DUPLICATE KEY UPDATE、Hive 的LATERAL VIEW、PostgreSQL 的-操作符。我一开始统一用readmysql解析结果 Hive 源头的 SQL 报废了 30% 以上。后来改成按来源标签选择方言解析失败率降到了 5% 以下。解析不了的 SQL 不要直接丢进入异常池定期人工看。3.2 基于 AST 的去重、难度分级和 Schema 绑定文本级去重在 SQL 场景下效果很差。SELECT * FROM users WHERE status 1和SELECT * FROM USERS WHERE STATUS 1文本不同但语义一模一样。所以去重必须拿到 AST 层面做签名比较。我用的签名方法大致是这样import hashlib def ast_signature(ast): normalized ast.copy() # 把所有字符串和数字字面量替换为占位符 for literal in normalized.find_all(sqlglot.exp.Literal): literal.replace(sqlglot.exp.Literal.string(?)) # 标识符统一小写 for identifier in normalized.find_all(sqlglot.exp.Identifier): identifier.set(this, identifier.name.lower()) return hashlib.sha256(normalized.sql().encode()).hexdigest()这样处理之后WHERE id 1和WHERE id 2会被认为是同一条模板保留一条、丢弃其他。当然这种强去重会损失部分数值多样性所以我做了一个折中配置模板级去重保留一条但允许同一个模板最多保留三条带不同数值的样本这样既控制了重复度又保留了条件取值的多样性。难度分级我采用打分制。基础分 1 分每个 JOIN 加 1 分GROUP BY 加 1 分HAVING 加 1 分每个子查询加 2 分每个 CTE 加 2 分每个窗口函数加 2 分WHERE 条件数量超过 3 个加 1 分嵌套层级超过 2 层加 3 分。总分 1-2 分是简单、3-5 分是中等、6 分以上是困难。这个打分规则不用非常精确关键是能保证训练集里各难度档位有足够样本避免模型只会答简单题。Schema 绑定是另一个重要环节。每条 SQL 都得对应一份建表 DDL模型才知道查到的是什么表。我的做法是解析 SQL 涉及的库表之后从元数据中心拉取真实 DDL然后做一遍字段裁剪——只保留当前 SQL 用到的字段其余字段用注释说明“存在但未使用”。这样既能降低模型的上下文长度压力又能让模型关注当前查询真正需要的字段。裁剪后的 DDL 和 SQL 一起存进训练样本。3.3 执行验证让每条 SQL 都是“能跑的”文本上合法、但在真实库上报错的 SQL 是 Text-to-SQL 模型最常见的“幻觉”来源。训练数据里绝对不能出现这种坏样本。执行验证我分两层做。第一层是 EXPLAINimport pymysql conn pymysql.connect( hostshadow_db_host, port3306, uservalidator, password****, databaseshadow_schema ) with conn.cursor() as cur: try: cur.execute(fEXPLAIN {sql}) rows cur.fetchall() if not rows: raise ValueError(EXPLAIN 无返回疑似非查询语句) except Exception as e: # 验证失败SQL 进入“不可执行池” print(f验证失败: {e}) finally: conn.close()EXPLAIN 能挡掉大部分语法错误和表不存在问题。但 EXPLAIN 只能证明“能解析”不能证明“能跑出合理结果”所以还有第二层抽样执行。我在影子库上对全表做TABLESAMPLE之后给 SQL 强制套一层 LIMIT 1SELECT * FROM ({sql}) AS _valid_check LIMIT 1这样既验证了真实执行路径又不会因为查询量太大压垮影子库。如果这一步能把结果跑出来基本可以判定这条 SQL 是“活”的。执行耗时超过 5 秒的我也直接过滤训练数据里的 SQL 应该是有效率的写法让模型学到低效查询绝对是个灾难。这里有一个容易踩的坑MySQL 的 EXPLAIN 并不执行 SQL但 EXPLAIN ANALYZE 会。如果你用的是 EXPLAIN ANALYZE等于直接跑了一遍全量查询在线上日志量大的时候这个开销完全扛不住。验证阶段一定要区分清楚。3.4 自然语言问题的生成与双向往返校验有了 SQL 还不够训练样本还需要对应的自然语言问题。理想情况下线上 BI 报表名称、数据需求文档里都有现成的自然语言描述但实际情况是大部分 SQL 根本没有配套文案需要额外生成。我的方案是用一个大模型做“SQL→NL”翻译再用另一个模型做“NL→SQL”反译最后比较两条 SQL 的语义一致性。这本质上是把一致性问题变成了可验证问题# 伪代码示意 sql_A SELECT dept_id, COUNT(*) FROM employees WHERE hire_date 2023-01-01 GROUP BY dept_id ORDER BY COUNT(*) DESC LIMIT 5 # 第一步SQL - 自然语言 nl llm_generate( 把下面这条 SQL 翻译成数据分析师能理解的自然语言问题不要包含表名和字段名\n sql_A ) # 期望输出展示2023年之后入职人数最多的5个部门 # 第二步自然语言 - SQL sql_B llm_generate( 根据下面的数据库 Schema 和自然语言问题写出对应的 SQL\n schema_ddl \n问题 nl ) # 第三步语义一致性比较 if compare_sql_semantics(sql_A, sql_B): accepted.append({schema: schema_ddl, question: nl, sql: sql_A})compare_sql_semantics的判定我采用“执行结果对比 AST 结构对比”双重策略两个 SQL 在影子库上的执行结果集合一致并且 AST 关键算子一致才判定为通过。结果对比必须注意 ORDER BY 的影响集合对比要忽略行序只比较行集合的哈希。双向往返校验的通过率是个很好的质量指标。我跑了几轮之后发现SQL 越复杂往返一致率越低复杂嵌套查询的往返一致率只有 60% 左右。这个阶段筛掉的样本不要扔掉它们是很好的“困难样本”可以专门攒起来后续让更专业的标注人员手工修正后加入训练集。4. DataFlow 管道实践节点配置与调度4.1 定义数据资源和引擎DataFlow 的编排思路不适合用脚本一把梭哈去写。它有点像一个数据版本的“DAG 控制系统”你需要先声明每个节点的输入输出到底是什么。先定义资源。我把 MySQL 审计日志、Git 仓库 SQL 文件、数据仓库查询历史分别声明成输入资源把训练数据集存储声明为输出资源。Resource 上要配置好读取协议、数据格式、刷新策略节点不直接感知资源细节。然后是引擎。开发阶段我用本地线程池引擎直接跑单机测试数据量小、调试方便全量跑批的时候切到 Flink 引擎把节点函数打包提交到集群。切换引擎不需要改节点代码DataFlow 的运行时层把资源映射和调度策略打包处理了。这是这个设计最有价值的点代码不改环境一跳。我经常在笔记本上跑 1 万条数据做逻辑验证然后直接发布到生产跑 1000 万条全量跑法不同、逻辑绝对一致。4.2 核心节点的参数与策略配置整个 Pipeline 我用 YAML 清单来定义方便版本管理和配置修改pipeline: name: text2sql_data_pipeline engine: flink nodes: - name: sql_collector type: source resource: mysql_audit_log params: time_window: 2024-01-01 00:00:00 ~ 2024-06-30 23:59:59 - name: sql_parser type: transform function: parse_sql_to_ast params: dialect_map: mysql: mysql hive: hive presto: presto - name: sql_filter type: transform function: filter_valid_sql params: allow_types: [select] max_length: 2000 max_nesting_depth: 5 - name: sql_dedup type: transform function: dedup_by_ast_signature params: template_keep_limit: 3 - name: sql_validator type: transform function: validate_execution params: shadow_db_host: 10.20.30.40 execute_timeout_ms: 5000 - name: sample_generator type: transform function: generate_training_sample params: llm_model: qwen2.5-72b-instruct round_trip_check: true - name: dataset_writer type: sink resource: oss_training_data params: format: jsonl注意几个参数的设置逻辑。template_keep_limit: 3前面解释过是为了在去重和多样性之间折中。time_window用来控制数据回溯范围我第一次跑的时候就吃了这个亏没有设置窗口把一整年的审计日志全拉进来了结果 80% 的数据是历史 obsolete 表验证阶段大量失败白白浪费了 6 个小时的集群资源。execute_timeout_ms: 5000是为了保护影子库线上日志里会有一些故意写得很慢的探测查询让验证阶段快速失败比让它们跑几十秒更明智。4.3 增量更新与失败重试的细节整条 Pipeline 跑全量看起来挺美好但业务库每天都在变SQL 日志持续产生只用一次性全量处理很快就会过时。我在实践中把管线拆成两种节奏全量任务每天凌晨跑一次负责把当天新增日志、新增表结构变化融进来。增量任务每隔 10 分钟跑一次只处理新到的日志片段。增量处理的关键是位点管理——记录每条日志的来源位置和时间戳处理完之后把位点更新到状态表里。下次增量任务从上次的位置继续而不是从零开始扫日志。DataFlow 对资源上的状态位点有内建支持不需要每份数据自己维护标记。失败重试这个环节我在真实环境里踩过不少坑关键是每个处理节点必须做到幂等。所谓幂等就是同一份输入被重复处理两次结果没有任何不同。我把最终输出设计成主键可覆盖写每条训练样本用“SQL 签名 指纹 版本号”作为主键重复写入直接覆盖旧记录。这样即使某个节点挂掉导致重跑前面几个节点最终写到训练集里的数据不会出现重复。上游节点可以随便重试下游节点靠幂等兜底。还有一个我没预料到的细节验证阶段的并发控制。全量跑的时候验证节点同时发起大量 EXPLAIN 请求影子库直接被连接数打满了。后来加了信号量并发限制把同一时刻执行验证的 SQL 数量控制在 20 以内影子库才稳定住。不管引擎多强下游数据库的资源永远可能成为瓶颈这种并发控制参数一定要在节点配置里提前设计好。5. 常见问题与排查实录5.1 解析失败与方言兼容解析阶段报错是最常见的几乎每隔几天就能筛出一批新花样。有一次我发现线上 MySQL 的日志里出现了SELECT SQL_NO_CACHE开头的语句sqlglot 默认方言在某些版本下解析不了这个语法整批报错。排查了很久才发现是监控工具自动打的标记。解决方案是在清洗规则里增加前缀裁剪把SQL_NO_CACHE、SQL_CALC_FOUND_ROWS这类调试性关键字在解析前剥离。这类问题的好处是规则积累起来之后解析失败率会越降越低。另外一个实战经验不要把解析失败率当成一个可以假装不存在的指标。我在 Pipeline 里给解析节点加了一个“异常池”输出所有解析失败的 SQL 单独落盘每周做一次原因归类。归类下来发现失败原因里方言语法问题占一半、截断文本占三成、编码问题占两成。每一类都有对应的修法改完之后整体通过率从 83% 提升到了 94%。5.2 Schema 漂移怎么处理这是所有 Text-to-SQL 数据管线最难的问题。线上库表结构经常变字段改名、表拆分成两张、某列类型从 INT 改成 BIGINT。SQL 日志里的历史查询可能引用的是三个月前的表结构直接拿到今天的影子库上验证大概率报“Unknown column”。我的处理策略是“按时间窗对齐 Schema”。从元数据中心拉取每一个时间窗口内的 DDL 快照验证 SQL 时用的是 SQL 发生时刻对应的 Schema而不是当前最新 Schema。实现方式是建立一张 Schema 版本表记录每个库表结构的生效起止时间SQL 先根据采集时间定位到对应的 Schema 版本再决定怎么验证。如果实在找不到历史 Schema 版本那就直接降级处理只要 SQL 在最新 Schema 上能跑通就保留跑不通就进“待确认池”交给人工判断。在实践中这个待确认池大约占总量的 5-8%完全可以接受。一开始我不知道这个比例总觉得所有过期 SQL 都该保留结果验证通过率只有 70%优化 Schema 对齐策略之后直接拉到了 90% 以上。5.3 数据多样性不足的优化第一版 Pipeline 跑出来的训练集有个非常明显的问题单表简单查询占了 60% 以上多表 JOIN 和窗口函数查询只占不到 15%。模型在这种分布下训练完遇到需要三个表以上的 JOIN 时表现很差。这个问题的根因不在模型在采样逻辑。线上日志本来就是简单查询占多数不能期待自然分布刚好满足训练需求。我的解决办法是在采样节点加了“难度桶均衡”策略先按维度标签分桶再在每个桶内做配额采样。简单查询最多保留 20%中等难度保留 40%困难查询保留 40%。另外按照算子维度做二级约束比如窗口函数至少占 5%CTE 至少占 8%多层嵌套至少占 10%。加了均衡策略之后训练集的分布合理多了模型在验证集上的复杂查询准确率提升了差不多 15 个百分点。这个经验其实挺通用的训练数据分布要服务于任务目标而不是忠实反映原始数据分布。5.4 安全与隐私脱敏怎么做线上 SQL 直接进训练集涉及数据安全的问题比想象中多。SQL 里经常出现真实的手机号、身份证号、订单金额这些绝对不能出现在训练数据里。脱敏处理我采用“语义保持替换”策略识别出 PII 类型的字面量之后不删除而是替换成同类型的伪造数据。比如手机号正则匹配到的138******格式的字符串替换成13900001111这样的测试号码身份证号替换成格式合法的假号金额数字替换成接近的随机值。这样替换的优点是 SQL 的 WHERE 比较逻辑不变模型学到的模式依然有效。表名和字段名也要脱敏。线上表名经常叫t_oms_order_detail、tb_user_credit_v2直接暴露业务信息。我维护了一张“命名映射表”把真实表名映射成语义等价的通用名比如t_oms_order_detail映射成order_detail字段里的user_phone映射成contact_phone。映射表同时作用于 DDL 和 SQL 文本保证 Schema 里表名和 SQL 里的引用是配套的。还有一个容易被忽略的点自然语言问题里也可能泄露敏感信息。LLM 生成 NL 时如果原 SQL 的 WHERE 条件里有真实地名、门店名、客户名生成的问句也会带上。这个只能靠后置敏感词过滤我在 sample_generator 节点后面接了一个敏感词扫描命中直接整条丢弃。6. 一些额外的实践心得跑完这套管线我最大的感受是Text-to-SQL 模型的天花板很大程度上不取决于模型而取决于训练数据的“干净度 难度覆盖度 Schema 对齐度”。大模型微调对数据质量极其敏感坏样本的干扰力远远大于好样本的提升力一条语义不一致的数据可能需要二十条高质量数据才能抵消。具体到 DataFlow 这套编排方案有几个教训值得记住。第一不要试图把清洗逻辑全塞进一个大函数里节点拆得越细问题定位越容易重跑成本也越低。第二每个处理节点都要带统计指标输出比如“输入多少条 / 过滤掉多少条 / 过滤原因分布”否则出了问题根本无从查起。第三训练集构造完之后一定要留一份“可追溯映射”每条样本都能定位回原始 SQL、原始日志、处理版本不然后续调模型时想找某条样本是怎么来的都找不到。如果让我重新做一遍我会在一开始就把“执行结果快照”也存下来。也就是说每条训练样本除了 Schema、NL、SQL 之外再保存一份它在影子库上的真实执行结果样例。这样后续做 RLHF 时可以用执行结果作为奖励信号让模型学会“写成这样能跑出合理结果”。这个扩展不需要改管线架构只需要在 validator 节点多输出一列结果快照。把这个设计提前考虑进去能省掉后面一次大返工。