大模型辅助 ETL 开发:从需求描述到生产级 Pipeline 与 dbt 模型生成
大模型辅助 ETL 开发:从需求描述到生产级 Pipeline 与 dbt 模型生成
在大多数数据团队中,工程师日常至少有 40% 的时间消耗在机械性的“提数与基础 ETL 拼装”上:业务方提来一段长篇大论的需求,诸如“把近 30 天华东和华南的已完成订单按天汇总,排除测试商户与退款单,关联用户维表计算新客占比,并写入 DWS 汇总表”。
资深工程师写这类 SQL 或 dbt 模型轻车熟路,但过程极其枯燥:查元数据找表名、拼 JOIN 条件、写 CTE 临时表、核对分区字段、编写 dbt 测试 YAML。一旦需求堆积,数据响应周期就会从“小时级”拖慢到“周级”。
大语言模型(LLM)的出现为 ETL 开发提供了全新的可能。但从“自然语言 Demo”走向“企业级生产流水线”,最怕的就是让模型直接裸写 SQL 并合入主干——模型可能会捏造字段、写出笛卡尔积 JOIN,甚至在没有增量策略的情况下直接全表覆盖。
真正稳妥的企业级方案,必须采用规范驱动生成(Spec-Driven Synthesis)与确定性编译验证(Deterministic Verification)。本文深入剖析基于 LLM 自动生成标准化 dbt 模型与质量测试的架构与工程实现。
一、为什么直接自然语言生成 ETL 容易翻车?
在大模型处理复杂 ETL 管道时,主要面临四大技术瓶颈:
+-----------------------------------------------------------------------------------+
| 痛点 1: 粒度失配与扇出膨胀 (Fan-out Trap) |
| 订单事实表 (1 对多) 关联物流轨迹表时,若模型未做去重直接 JOIN,会导致交易额翻倍。 |
+-----------------------------------------------------------------------------------+
| 痛点 2: 缺乏企业级工程规范 (CTE 与命名规范) |
| 生成出上百行多层嵌套子查询的“意大利面条代码”,后期根本无法维护与排查。 |
+-----------------------------------------------------------------------------------+
| 痛点 3: 缺少增量物化与幂等设计 (Incremental Strategy) |
| 默认写成全表全量扫描,无法自动适配基于 `is_incremental()` 的水印合并逻辑。 |
+-----------------------------------------------------------------------------------+
| 痛点 4: 缺少数据契约与伴生测试 (Data Contract & Tests) |
| 仅生成了 SQL 代码,缺少 `schema.yml` 中的非空、唯一性与外键约束声明。 |
+-----------------------------------------------------------------------------------+
因此,必须把大模型的职责严格收敛为:“根据元数据上下文,生成符合企业规范的 CTE 结构、Jinja 引用与伴生测试用例”,后续由静态分析器(如 SQLGlot、SQLFluff)与 CI/CD 流水线完成确定性检查。
二、端到端生成架构:规范驱动的 ETL 合成链路
一条严谨的 ETL 自动合成流水线分为四个核心环节:
+-----------------------------------------------------------------------------------+
| 1. 上下文与元数据注入 (Metadata RAG) |
| - 检索源表的 DDL、字段注释、主外键拓扑与历史优质 dbt 模型模板 |
+-----------------------------------------------------------------------------------+
|
v
+-----------------------------------------------------------------------------------+
| 2. 结构化提示词与规范约束 (Prompt with CTE Guidelines) |
| - 强制采用标准 CTE 架构: import_ctes -> logical_ctes -> final_select |
| - 强制使用 dbt Jinja: {{ source() }} 与 {{ ref() }},禁止硬编码物理表名 |
+-----------------------------------------------------------------------------------+
|
v
+-----------------------------------------------------------------------------------+
| 3. 代码与伴生契约联合生成 (Co-Generation) |
| - 产物 A: `models/dws_trade_daily_summary.sql` (增量模型代码) |
| - 产物 B: `models/dws_trade_daily_summary.yml` (字段契约与 Great Expectations) |
+-----------------------------------------------------------------------------------+
|
v
+-----------------------------------------------------------------------------------+
| 4. 静态编译与沙箱 Dry-Run 校验 (Verification Pipeline) |
| - SQLGlot 语法树检查: 验证字段存在性、禁止破坏性 DDL |
| - dbt compile & dbt test --dry-run: 验证 Jinja 解析与执行计划 |
+-----------------------------------------------------------------------------------+
三、生产级 dbt 规范:现代数仓的 CTE 标准结构
为了保证大模型生成的代码具备工业级可读性与可维护性,团队通常要求模型遵循 dbt Labs 推荐的 CTE 编码范式:
{{
config(
materialized='incremental',
unique_key=['dt', 'city_id'],
partition_by={'field': 'dt', 'data_type': 'date'},
incremental_strategy='merge'
)
}}
-- 1. Import CTEs: 仅声明源表引用与基础增量过滤
WITH source_orders AS (
SELECT *
FROM {{ ref('fct_orders') }}
{% if is_incremental() %}
WHERE updated_at >= (SELECT MAX(updated_at) FROM {{ this }})
{% endif %}
),
dim_users AS (
SELECT *
FROM {{ ref('dim_users') }}
),
-- 2. Logical CTEs: 处理业务过滤、去重与指标计算
valid_orders AS (
SELECT
order_id,
user_id,
city_id,
pay_amount,
DATE(created_at) AS dt
FROM source_orders
WHERE order_status = 'COMPLETED'
AND is_test_order = 0
),
joined_details AS (
SELECT
o.dt,
o.city_id,
o.order_id,
o.pay_amount,
u.is_new_buyer
FROM valid_orders o
LEFT JOIN dim_users u ON o.user_id = u.user_id
),
-- 3. Final Select: 最终聚合与指标命名
final AS (
SELECT
dt,
city_id,
COUNT(DISTINCT order_id) AS total_order_cnt,
SUM(pay_amount) AS total_pay_amt,
COUNT(DISTINCT CASE WHEN is_new_buyer = 1 THEN order_id END) AS new_buyer_order_cnt
FROM joined_details
GROUP BY dt, city_id
)
SELECT * FROM final
四、生产级生成引擎与 AST 静态校验器实现
下面的 Python 实现结合了 sqlglot 语法分析器,能够自动化验证模型生成的 SQL 是否存在非法表名、是否包含未授权函数、是否提取了血缘依赖,并自动构建对应的 schema.yml 数据测试文件。
"""
llm_etl_compiler.py
生产级大模型辅助 ETL 与 dbt 模型生成与安全编译引擎
"""
import json
from dataclasses import dataclass, field
from typing import Any, Dict, List, Optional, Set
import sqlglot
from sqlglot import exp
import yaml
@dataclass
class ColumnMeta:
name: str
data_type: str
description: str
is_primary_key: bool = False
tests: List[str] = field(default_factory=list)
@dataclass
class TableMetadata:
table_name: str
columns: Dict[str, ColumnMeta]
class ETLVerificationError(Exception):
pass
class DbtModelCompiler:
"""
负责模型生成的 SQL 静态 AST 分析、字段血缘校验与测试契约构建
"""
def __init__(self, allowed_tables: Dict[str, TableMetadata]):
self.allowed_tables = allowed_tables
def verify_and_analyze_sql(self, sql_code: str, target_dialect: str = "spark") -> Dict[str, Any]:
"""
利用 SQLGlot 解析 SQL AST,严格检查:
1. 语法正确性与方言适配
2. 禁止危险 DDL/DML 操作 (DROP, TRUNCATE, DELETE)
3. 提取所引用的源表与字段血缘
"""
try:
# 剥离 Jinja 宏标签进行 AST 纯语法解析
clean_sql = sqlglot.transpile(sql_code, read=target_dialect, write=target_dialect)[0]
ast = sqlglot.parse_one(clean_sql, read=target_dialect)
except Exception as e:
raise ETLVerificationError(f"SQL 语法解析失败: {e}")
# 检查是否包含危险操作
for node in ast.walk():
if isinstance(node, (exp.Drop, exp.Delete, exp.Truncate)):
raise ETLVerificationError(f"安全拦截: 禁止在 ETL Pipeline 中执行破坏性操作 {type(node).__name__}")
# 提取表依赖血缘
tables_referenced = set()
for table_expr in ast.find_all(exp.Table):
tables_referenced.add(table_expr.name)
# 提取输出字段列表
output_columns = []
for select in ast.find_all(exp.Select):
for expr in select.expressions:
if isinstance(expr, exp.Alias):
output_columns.append(expr.alias)
elif isinstance(expr, exp.Column):
output_columns.append(expr.name)
return {
"is_valid": True,
"tables_referenced": list(tables_referenced),
"output_columns": list(set(output_columns)),
}
def generate_dbt_schema_yaml(self, model_name: str, description: str, columns: List[ColumnMeta]) -> str:
"""
自动生成配套的 dbt schema.yml 测试契约文件
"""
cols_payload = []
for col in columns:
col_dict = {
"name": col.name,
"description": col.description,
}
tests = list(col.tests)
if col.is_primary_key:
if "unique" not in tests:
tests.append("unique")
if "not_null" not in tests:
tests.append("not_null")
if tests:
col_dict["tests"] = tests
cols_payload.append(col_dict)
schema_dict = {
"version": 2,
"models": [
{
"name": model_name,
"description": description,
"columns": cols_payload
}
]
}
return yaml.dump(schema_dict, sort_keys=False, allow_unicode=True)
运行验证与契约生成测试
# 1. 模拟数仓元数据
METADATA_REPO = {
"fct_orders": TableMetadata(
table_name="fct_orders",
columns={
"order_id": ColumnMeta("order_id", "string", "订单主键", is_primary_key=True),
"pay_amount": ColumnMeta("pay_amount", "decimal(18,2)", "支付金额"),
"city_id": ColumnMeta("city_id", "string", "城市编码"),
"dt": ColumnMeta("dt", "string", "分区日期")
}
)
}
# 2. 模拟 LLM 生成的规范 dbt 模型 SQL
MOCK_GENERATED_SQL = """
SELECT
dt,
city_id,
COUNT(DISTINCT order_id) AS order_cnt,
SUM(pay_amount) AS gmv
FROM fct_orders
WHERE dt = '2026-08-24'
GROUP BY dt, city_id
"""
# 3. 执行 AST 校验与依赖分析
compiler = DbtModelCompiler(METADATA_REPO)
analysis = compiler.verify_and_analyze_sql(MOCK_GENERATED_SQL, target_dialect="spark")
print("AST 语法与依赖分析结果:\n", json.dumps(analysis, indent=2, ensure_ascii=False))
# 4. 生成配套的 dbt schema.yml 质量契约
model_columns = [
ColumnMeta("dt", "string", "统计日期", tests=["not_null"]),
ColumnMeta("city_id", "string", "城市唯一编码", tests=["not_null"]),
ColumnMeta("order_cnt", "bigint", "总完单数", tests=["not_null"]),
ColumnMeta("gmv", "decimal(18,2)", "总成交金额", tests=["not_null"])
]
schema_yml = compiler.generate_dbt_schema_yaml(
model_name="dws_trade_city_daily",
description="分城市每日成交与下单汇总模型",
columns=model_columns
)
print("\n生成的配套 schema.yml 契约:\n" + schema_yml)
五、企业级落地工程防线与协同工作流
为了确保大模型生成的 ETL 代码不会污染生产数仓,必须在 CI/CD 中设置三层防护网:
- 自动提 PR 并附带数据 Diff 报告:
大模型生成代码后,严禁直接 Merge 到主分支。系统会自动创建 Git 分支并在影子数仓中试跑,比对新模型与旧模型的产出指标差值(Data Diff),生成可视化核对报告贴在 PR 评论区。 - 强制 Lint 与代码风格审查(SQLFluff):
强制检查关键字大小写、缩进、逗号前置/后置规范以及别名命名风格,确保团队代码风格高度一致。 - 核心财务模型双人 Review 制:
对于标有tier_1(核心财报与高管大屏)标签的数据模型,必须保留两位资深数据开发工程师的线下签字审批。
大模型不是要替代数据工程师,而是将工程师从繁重的模板语法与基础拼装中解脱出来,把精力集中在真正关键的业务指标定义、数仓维度建模与架构性能优化上。
更多推荐


所有评论(0)