大模型辅助 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 中设置三层防护网:

  1. 自动提 PR 并附带数据 Diff 报告
    大模型生成代码后,严禁直接 Merge 到主分支。系统会自动创建 Git 分支并在影子数仓中试跑,比对新模型与旧模型的产出指标差值(Data Diff),生成可视化核对报告贴在 PR 评论区。
  2. 强制 Lint 与代码风格审查(SQLFluff)
    强制检查关键字大小写、缩进、逗号前置/后置规范以及别名命名风格,确保团队代码风格高度一致。
  3. 核心财务模型双人 Review 制
    对于标有 tier_1(核心财报与高管大屏)标签的数据模型,必须保留两位资深数据开发工程师的线下签字审批。

大模型不是要替代数据工程师,而是将工程师从繁重的模板语法与基础拼装中解脱出来,把精力集中在真正关键的业务指标定义、数仓维度建模与架构性能优化上。

Logo

有“AI”的1024 = 2048,欢迎大家加入2048 AI社区

更多推荐