十年匠心定制 · 商业建站与技术教学双线并行 咨询热线:400-886-1026 service@lmnt.cn
ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

自然语言驱动ETL:从SQL取数到智能数据管道实战

自然语言驱动ETL:从SQL取数到智能数据管道实战 最近在整理数据中台需求时经常遇到一类场景业务同学想分析订单数据但每次都要提工单、等排期最终拿到的还是一张 SQL 结果表。如果能把“帮我看看 2024 年 6 月之后北京用户下单金额最高的前 3 个商品分类”这种自然语言请求直接翻译成一条可执行的 ETL 管道取数链路会轻松很多。TamedTable 这个项目从命名和定位来看就是把“以自然语言驱动 ETL”这件事产品化。本文会从 ETL 的基本概念讲起拆解自然语言转 ETL 的核心原理再用一个可运行的 Python 实战案例演示整套流程最后补充落地过程中的常见坑点与工程建议。1. 背景与核心概念1.1 什么是 ETLETL 是 Extract、Transform、Load 的缩写分别对应数据抽取、数据转换和数据加载。举个例子业务库中的订单表可能字段命名混乱、状态码含义不统一、时间格式五花八门此时需要把多份数据源抽取出来经过清洗、过滤、聚合、关联等转换操作再加载到数据仓库或分析系统中。ETL 是数据分析和数据仓库建设的基础环节。没有 ETL数据就是一堆散落的原始日志和业务表很难直接支撑报表、BI 看板和算法模型。传统 ETL 开发通常由数据工程师编写 SQL 或使用 DataX、Kettle、Airflow、dbt 等工具完成整个过程依赖人工理解业务需求再把需求翻译成技术实现。1.2 传统 ETL 的痛点传统 ETL 开发有几个非常明显的痛点第一沟通成本高。业务同学的描述往往是“我要看最近一个月各地区的销售情况”但这条需求落到数据工程师面前需要继续追问“按哪个粒度统计”“金额口径是含税还是不含税”“地区字段是用户所在地还是收货地”。一来一回半天就过去了。第二取数效率低。临时取数需求占了数据团队日常工作的很大比例。每个需求都走工单流转开发一条临时 SQL跑完再回传结果链路长、响应慢。第三数据管道难以复用。很多取数逻辑其实大同小异但因为每次都是临时写 SQL过滤条件、聚合口径、输出格式都散落在一个个脚本里没有沉淀成可复用的管道。第四对非技术人员不友好。业务同学有数据需求时要么学着写 SQL要么依赖数据团队两者都不是最优解。1.3 TamedTable 是什么TamedTable 是一个以自然语言为核心交互方式的 AI ETL 项目。简单来说用户不再需要关心 SQL 怎么写、管道怎么配只需要用中文描述“我想要什么样的数据”系统会自动理解语义生成对应的抽取、转换、加载逻辑并执行。这个思路可以拆成两层交互层用自然语言替代 SQL 和配置项。执行层底层仍然是结构化的 ETL 管道只是管道由 AI 自动生成。从工程角度看TamedTable 并不是要取代成熟的 ETL 调度平台而是把“自然语言 - 管道定义 - 执行结果”这条链路打通。它的核心价值在于降低数据消费门槛让非技术人员也能自助取数。1.4 适用场景与边界自然语言 ETL 比较适合以下场景内部数据问答业务同学用中文问数据系统返回统计结果。临时取数自动化将高频临时需求沉淀为可复用的自然语言指令。数据报表前置探索在生成固定报表之前先用自然语言探索数据口径。数据质量检查用自然语言描述“找出金额为负的订单”快速定位异常数据。但它也有明显的边界。自然语言 ETL 并不适合完全替代专业数据工程师的复杂管道开发。比如涉及多表复杂关联、自定义 UDF、窗口函数、增量更新策略时AI 生成的结果可靠性会下降仍然需要人工校验和干预。2. 环境准备与版本说明2.1 运行环境本文中的实战案例以 Python 3.10 为例操作系统不限Linux、macOS、Windows 都可以运行。代码主要使用 Python 标准库不需要额外安装重量级框架。版本需要根据你的项目实际情况调整下面的示例以常见环境为例重点演示实现思路。2.2 依赖库说明核心依赖只有两个pandas用于简化数据处理非必需但很方便。openai 或 requests用于调用大模型接口本文会给出两种思路。如果你只需要跑通本地规则解析版本不安装任何第三方库也可以。完整示例会尽量保持轻量方便读者直接复制运行。2.3 项目目录结构建议按下面的结构组织项目tamedtable-demo/ ├── data.py # 模拟订单数据 ├── field_map.py # 字段别名映射与 Schema 定义 ├── parser.py # 自然语言解析逻辑 ├── executor.py # ETL 管道执行器 ├── llm_parser.py # 大模型解析接入示例 └── main.py # 主流程入口模块之间的依赖关系是main.py 调用 parser.py 或 llm_parser.py得到统一的管道描述executor.py 负责执行管道data.py 提供测试数据。3. 核心原理拆解3.1 自然语言转 ETL 的四个阶段把一句自然语言变成可执行的 ETL 管道本质上是一个“语义到结构”的转换问题。整个过程可以拆成四个阶段意图识别判断用户想要做什么操作比如过滤、排序、分组、聚合还是合并。字段映射把自然语言中的业务名词映射到实际表字段例如“地区”对应 region 字段。条件解析抽取过滤条件例如“2024 年 6 月之后”对应 order_date 2024-06-01。管道生成把前面解析结果组装成一个有序的执行步骤列表。这四个阶段不一定都要用大模型完成。对字段映射和条件解析这类确定性较强的部分可以先用规则处理只有规则处理不了再交给大模型。3.2 意图识别把请求映射为管道步骤ETL 管道可以抽象成一系列操作步骤常见的操作包括操作类型动作关键词管道描述中的字段filter过滤、筛选、排除、只保留field, operator, valuegroup_by按…分组、分类统计keyaggregate求和、平均、最大、最小、计数func, field, aliassort排序、最高的、最低的sort_key, orderlimit前 N 条、Top Nlimit意图识别要解决的就是从“按商品分类统计总金额金额最高的前 3 条”中提取出 group_by、aggregate、sort、limit 四个步骤。3.3 字段映射与条件解析字段映射是自然语言 ETL 中最容易出错的一环。用户很少会说出精确的字段名更多是使用业务别名。例如“地区” - region“金额” - amount“订单日期” - order_date“商品分类” - category条件解析同样需要做语义映射。比如“6 月之后”要理解成order_date 2024-06-01而不是把字符串“6 月之后”直接放进 SQL。更稳妥的做法是定义一个字段元数据表把字段名、类型、中文别名、允许的操作符、取值范围全部维护起来。这不仅能提高解析准确率也是后续做安全校验的基础。3.4 生成可执行的 ETL 流水线解析完成后需要把结果输出成统一的管道描述。推荐使用 JSON 格式方便后续扩展和可视化。下面是一个示例[ {op: filter, conditions: [{field: order_date, operator: , value: 2024-06-01}, {field: region, operator: , value: 北京}]}, {op: group_by, key: category}, {op: aggregate, func: sum, field: amount, alias: total_amount}, {op: sort, sort_key: total_amount, order: desc}, {op: limit, limit: 3} ]得到这样的管道描述之后执行器只需要按顺序逐条执行即可。这种方式也方便记录日志和审计因为每一步都是结构化数据。3.5 为什么不能直接让大模型写代码就跑有些读者可能会想既然大模型这么强为什么不直接让大模型生成 Python 代码或 SQL 执行呢原因有三点第一不可控。大模型生成的代码可能使用不存在的函数、错误的库名甚至带有潜在安全风险。第二不可审计。直接执行生成的 SQL 或代码无法告诉用户“这条结果是怎么算出来的”。一旦结果错误很难定位问题。第三不安全。如果允许 AI 生成的 SQL 直接连生产库执行会出现危险操作比如全表删除、跨表笛卡尔积、超大结果集导出。所以更推荐的方案是让大模型输出结构化的 JSON 管道描述而不是直接输出可执行代码。系统对 JSON 做严格校验后再由自己的执行器执行。这正是 TamedTable 这类工具应该走的路。4. 完整实战案例这一节会实现一个简化版的自然语言 ETL。先用纯规则解析跑通流程再提供接入大模型的扩展思路。4.1 准备模拟订单数据为了演示方便不使用外部数据库。直接在 Python 中定义一个订单列表字段包含订单号、用户名、地区、商品分类、金额、订单日期。# data.py orders [ {order_id: A001, user_name: 张三, region: 北京, category: 数码, amount: 2999, order_date: 2024-05-12}, {order_id: A002, user_name: 李四, region: 上海, category: 服装, amount: 399, order_date: 2024-06-20}, {order_id: A003, user_name: 王五, region: 北京, category: 数码, amount: 4599, order_date: 2024-06-25}, {order_id: A004, user_name: 赵六, region: 广州, category: 食品, amount: 129, order_date: 2024-07-01}, {order_id: A005, user_name: 孙七, region: 北京, category: 服装, amount: 599, order_date: 2024-07-15}, {order_id: A006, user_name: 周八, region: 上海, category: 食品, amount: 259, order_date: 2024-07-18}, {order_id: A007, user_name: 吴九, region: 北京, category: 食品, amount: 89, order_date: 2024-08-02}, {order_id: A008, user_name: 郑十, region: 广州, category: 数码, amount: 1999, order_date: 2024-08-10}, ]4.2 定义字段映射与 Schema字段映射表是整个自然语言 ETL 的地基。# field_map.py FIELD_ALIAS { 订单号: order_id, 订单id: order_id, 用户名: user_name, 用户: user_name, 地区: region, 城市: region, 商品分类: category, 分类: category, 类别: category, 金额: amount, 价格: amount, 订单日期: order_date, 日期: order_date, 下单日期: order_date, } FIELD_TYPE { order_id: string, user_name: string, region: string, category: string, amount: number, order_date: date, }这里要注意字段映射表必须覆盖业务同学最常用的说法。不要指望大模型能“凭空知道”你表里的字段是 region 还是 district只有把 Schema 信息明确告诉 AI解析准确率才会高。4.3 编写条件解析与管道生成下面的 parser.py 实现了一个轻量级的规则解析器。它通过关键词匹配把一条自然语言转换成管道描述。# parser.py import re from field_map import FIELD_ALIAS def extract_field(text): 从文本中提取字段名返回标准字段名。 for alias, field in FIELD_ALIAS.items(): if alias in text: return field return None def extract_region_condition(text): 识别地区条件例如北京用户 - region 北京。 for region in [北京, 上海, 广州, 深圳]: if region in text: return {field: region, operator: , value: region} return None def extract_date_condition(text): 识别日期条件例如2024年6月之后 - order_date 2024-06-01。 match re.search(r(\d{4})年(\d{1,2})月之后, text) if match: year, month match.group(1), match.group(2) date_value f{year}-{int(month):02d}-01 return {field: order_date, operator: , value: date_value} match re.search(r(\d{4})年(\d{1,2})月之前, text) if match: year, month match.group(1), match.group(2) date_value f{year}-{int(month):02d}-01 return {field: order_date, operator: , value: date_value} return None def parse_nl_to_pipeline(text): 自然语言 - 管道描述。 pipeline [] # 1. 过滤条件 conditions [] region_cond extract_region_condition(text) if region_cond: conditions.append(region_cond) date_cond extract_date_condition(text) if date_cond: conditions.append(date_cond) if conditions: pipeline.append({op: filter, conditions: conditions}) # 2. 分组条件按某字段分组 group_field None group_match re.search(r按(.?)统计|按(.?)分组|各(.?)的, text) if group_match: group_name group_match.group(1) or group_match.group(2) or group_match.group(3) group_field extract_field(group_name 的) if group_field: pipeline.append({op: group_by, key: group_field}) # 3. 聚合条件求和、平均、最大 agg_func None agg_field amount if 总金额 in text or 求和 in text or 总额 in text: agg_func sum elif 平均 in text or 均值 in text: agg_func mean elif 最大 in text: agg_func max elif 最小 in text: agg_func min if agg_func: # 如果文本中明确说了按哪个字段聚合优先使用该字段 for alias, field in FIELD_ALIAS.items(): if alias in text and alias ! 金额: agg_field field break pipeline.append({op: aggregate, func: agg_func, field: agg_field, alias: f{agg_func}_{agg_field}}) # 4. 排序最高的前 N 条 sort_field None if 最高 in text or 最多 in text or 最大 in text: sort_field agg_field elif 最低 in text or 最少 in text or 最小 in text: sort_field agg_field if sort_field: order desc if (最高 in text or 最多 in text or 最大 in text) else asc pipeline.append({op: sort, sort_key: sort_field, order: order}) # 5. 限制条数前 N 条 limit_match re.search(r前(\d), text) if limit_match: pipeline.append({op: limit, limit: int(limit_match.group(1))}) return pipeline上面的解析器虽然简单但已经能处理“2024年6月之后北京用户订单中按商品分类统计总金额金额最高的前3条”这类复合指令。4.4 编写 ETL 执行器执行器的核心思路很简单依次执行管道中的每个步骤。# executor.py from collections import defaultdict def execute_pipeline(rows, pipeline): result rows for step in pipeline: op step[op] if op filter: result _apply_filter(result, step[conditions]) elif op group_by: result _apply_group_by(result, step[key]) elif op aggregate: result _apply_aggregate(result, step[func], step[field], step[alias]) elif op sort: result _apply_sort(result, step[sort_key], step[order]) elif op limit: result result[: step[limit]] return result def _apply_filter(rows, conditions): filtered [] for row in rows: match True for cond in conditions: field, operator, value cond[field], cond[operator], cond[value] if operator and str(row.get(field)) ! str(value): match False break if operator and row.get(field) value: match False break if operator and row.get(field) value: match False break if match: filtered.append(row) return filtered def _apply_group_by(rows, key): groups defaultdict(list) for row in rows: groups[row.get(key)].append(row) result [] for group_key, group_rows in groups.items(): merged dict(group_rows[0]) merged[_group_rows] group_rows merged[_group_key] group_key result.append(merged) return result def _apply_aggregate(rows, func, field, alias): result [] # 如果已经按组分组rows 中每个元素是一个分组 if rows and _group_rows in rows[0]: for group in rows: group_rows group[_group_rows] values [r.get(field, 0) for r in group_rows] result.append({ group_key: group[_group_key], alias: _calc(func, values) }) else: values [r.get(field, 0) for r in rows] result.append({alias: _calc(func, values)}) return result def _apply_sort(rows, sort_key, order): reverse order desc return sorted(rows, keylambda x: x.get(sort_key, 0), reversereverse) def _calc(func, values): if func sum: return sum(values) if func mean: return round(sum(values) / len(values), 2) if func max: return max(values) if func min: return min(values) return values这里做了一点简化在_apply_group_by后每个分组会保留一个_group_rows字段聚合时再从分组中取出原始行计算。这个设计方便后续扩展比如同时计算总和与平均值。4.5 主流程入口最后把解析器和执行器串起来。# main.py from data import orders from parser import parse_nl_to_pipeline from executor import execute_pipeline def main(): text 2024年6月之后北京用户订单中按商品分类统计总金额金额最高的前3条 print(用户输入, text) print() pipeline parse_nl_to_pipeline(text) print(生成的管道描述) for step in pipeline: print(step) print() result execute_pipeline(orders, pipeline) print(执行结果) for row in result: print(row) if __name__ __main__: main()运行 main.py预期输出如下用户输入 2024年6月之后北京用户订单中按商品分类统计总金额金额最高的前3条 生成的管道描述 {op: filter, conditions: [{field: region, operator: , value: 北京}, {field: order_date, operator: , value: 2024-06-01}]} {op: group_by, key: category} {op: aggregate, func: sum, field: amount, alias: sum_amount} {op: sort, sort_key: sum_amount, order: desc} {op: limit, limit: 3} 执行结果 {group_key: 数码, sum_amount: 4599} {group_key: 服装, sum_amount: 599} {group_key: 食品, sum_amount: 89}手动验证一下2024 年 6 月之后北京用户有 A003数码 4599、A005服装 599、A007食品 89分组求和后按金额降序取前 3 条结果正确。4.6 接入大模型解析扩展规则解析器能覆盖简单场景但真实业务中的问法更复杂比如“对比一下北京和上海在数码品类的销售额差距”“帮我看看上月退货率最高的三个地区”。这些都需要大模型来理解语义。接入大模型时建议用结构化输出而不是让它直接生成 SQL。# llm_parser.py import json SYSTEM_PROMPT 你是一个 ETL 管道解析器。你会收到用户的数据分析请求需要输出 JSON 格式的管道描述。 管道描述必须使用以下结构 [ {op: filter, conditions: [{field: 字段名, operator: ||||, value: 值}]}, {op: group_by, key: 分组字段名}, {op: aggregate, func: sum|mean|max|min|count, field: 聚合字段名, alias: 输出字段名}, {op: sort, sort_key: 排序字段, order: asc|desc}, {op: limit, limit: 数字} ] 只输出 JSON不要输出任何解释。 可用字段如下 - order_id: 订单号 - user_name: 用户名 - region: 地区取值范围北京、上海、广州、深圳 - category: 商品分类 - amount: 金额数字类型 - order_date: 订单日期格式 YYYY-MM-DD def parse_with_llm(text, client, modelyour-model-name): 使用大模型解析自然语言。传入 OpenAI 兼容客户端。 response client.chat.completions.create( modelmodel, messages[ {role: system, content: SYSTEM_PROMPT}, {role: user, content: text}, ], temperature0.1, ) content response.choices[0].message.content.strip() # 这里建议增加 JSON 解析异常兜底 return json.loads(content)关键点在于提示词中必须显式给出字段列表和字段取值范围并要求只输出 JSON。模型返回后还要做白名单校验确保字段名真实存在、操作符合法、limit 不超过上限。校验逻辑可以复用 field_map.py 中的字段信息。大模型的参数名和调用方式可能因版本不同而变化示例代码需要按你实际使用的 SDK 调整。核心思想是“大模型只负责生成结构化管道不负责直接执行”。5. 常见问题与排查思路问题现象常见原因解决思路字段识别错误字段别名表覆盖不全或者提示词未给出字段信息扩充别名表把表结构和字段说明注入提示词大模型返回内容不是合法 JSON提示词约束不足或者模型本身输出不稳定在提示词中强调“只输出 JSON”使用函数调用能力增加 JSON 解析兜底逻辑条件解析错误例如“6 月之后”被解析成等于 6 月语义理解不到位引入日期特征词规则在提示词中给出示例输入输出聚合结果不对分组和聚合顺序错误或聚合字段取错检查管道步骤顺序打印中间结果对比手工计算结果数据量太大执行超时直接在内存中处理大量数据引入数据库或 DataFrame 引擎增加采样预览功能限制返回行数生成 SQL 直接在生产库执行但删除了数据未做操作白名单禁止 AI 直接生成可执行 SQL只允许结构化管道数据库账号使用最小只读权限6. 最佳实践与工程建议6.1 Schema 元数据是解析质量的基石自然语言 ETL 的解析效果很大程度上取决于你对模型的“喂养”质量。字段列表、字段类型、取值范围、枚举值、业务口径说明这些元数据越完整解析准确率越高。建议维护一份 YAML 或 JSON 格式的元数据文件包括字段名、类型、中文别名、示例值、备注信息。解析时把这份元数据注入提示词或用于规则匹配不要写死在代码逻辑里。6.2 提示词工程要点写提示词时有几个容易被忽略的点第一不要让模型猜字段。必须显式列出所有可用字段否则模型会用臆想的字段名生成管道。第二给出输入输出示例。示例比描述更容易让模型理解格式。尤其是日期处理、排序方向这类容易混淆的地方。第三采用低温度参数。解析任务对创造性要求低对准确性要求高temperature 建议设置在 0.1 到 0.3 之间。第四要求 JSON 输出。JSON 比自然语言可解析性更高也比直接生成代码更安全。6.3 安全边界白名单、降权、审计自然语言 ETL 本质上是把“数据提取能力”开放给了更多人必须做好安全控制。操作白名单只允许 filter、group_by、aggregate、sort、limit 等白名单操作不允许任意代码执行。字段白名单校验 JSON 管道中出现的字段名必须在元数据表中存在。最小权限底层数据库账号只分配只读权限禁止 DDL、DELETE、UPDATE。限流与超时限制单次查询的数据量和执行时间避免全表扫描拖垮数据库。审计日志记录每次自然语言输入、生成的管道、执行结果摘要方便追溯问题。6.4 结果可解释与数据血缘数据结果如果不解释口径业务同学很难放心使用。建议在返回结果时附带“管道说明”例如筛选条件地区 北京订单日期 2024-06-01 分组字段商品分类 聚合方式金额求和 排序方式按总金额降序 返回条数前 3 条这比直接抛出一个数字更友好也能让用户判断“这不是我想要的口径”从而及时调整问题表达方式。6.5 性能优化策略当数据量变大时建议做几件事第一把执行层从纯 Python 换成数据库或 DataFrame 引擎。pandas、DuckDB、ClickHouse 都能高效处理过滤、聚合和排序。第二提供数据预览模式。默认只返回前 100 或前 1000 条避免大结果集传输。第三开启结果缓存。相同或相似的自然语言请求可以在短时间内直接读取缓存结果。缓存 key 可以使用标准化后的管道 JSON。第四对常用查询做预计算。如果业务上每天都有人问“各地区的总销售额”可以提前跑好汇总表在线查询时直接读取预计算结果。6.6 生产落地路径从 Demo 到生产建议按以下节奏推进第一步先解决“有结果”。用规则解析器处理高频固定句式快速跑通流程。第二步接入大模型增强理解能力。覆盖多变的业务问法同时保留规则解析作为兜底。第三步增加反馈机制。对解析结果增加“满意/不满意”按钮不满意时记录用户原始输入定期分析错误模式再补充规则或优化提示词。第四步对接真实数据源。把 pipeline 执行器从内存数据替换为数据库查询逐步支持 Hive、ClickHouse、MySQL、PostgreSQL 等数据源。第五步完善权限、审计、监控体系再逐步推广给更多业务团队使用。7. 总结与学习路线这篇文章围绕 TamedTable 所代表的“AI ETL 自然语言”方向梳理了从 ETL 基础概念到自然语言管道生成、执行、排错、落地的完整思路。核心收获可以归纳为三点第一自然语言 ETL 的本质不是让 AI 代替你写 SQL而是让 AI 帮你把自然语言翻译成结构化的管道描述再由可靠引擎执行。第二字段映射和 Schema 元数据决定了解析质量的上限。大模型再强也需要业务上下文。第三步实践时可以先从规则解析搭建最小闭环再逐步接大模型、接真实数据源、完善安全审计。如果是对数据工程感兴趣可以进一步学习 dbt、Airflow、DuckDB理解真实生产环境中的管道编排和调度。如果是对大模型应用感兴趣可以深入研究函数调用、结构化输出、RAG 与数据问答的结合方式。建议你从本文的 demo 代码开始把数据换成自己业务里的表一边跑一边记录解析失败的问题不断扩充字段别名和提示词示例。这个过程本身就是一次完整的大模型工程实践。
返回列表