09. dbt 数据转换与 DataOps

dbt (data build tool) 完整实践指南:models、tests、docs、macros 核心概念,Jinja 模板编程,CI/CD 集成以及 dbt Cloud 与 dbt Core 部署策略。

dbt(data build tool)是数据工程领域最具影响力的开源工具之一,它将软件工程最佳实践(版本控制、模块化、测试、文档)引入数据转换流程。本文从 Core 概念到生产部署全面讲解 dbt 生态。

1. dbt 核心概念与架构

1.1 为什么用 dbt

传统 ETL 的痛点与 dbt 的解决方案:

痛点dbt 方案
SQL 散落在各种脚本/工具中统一的 SQL 项目结构
缺乏版本控制Git-native 工作流
数据质量靠人工检查内置测试框架
表血缘关系黑盒自动生成 Lineage DAG
文档与代码分离docs-as-code
环境隔离困难多环境配置(dev/staging/prod)

1.2 dbt 架构

┌─────────────────────────────────────────────────────────────┐
│                      dbt Project                             │
│  ┌────────┐ ┌────────┐ ┌────────┐ ┌────────┐ ┌────────┐  │
│  │ models │ │ tests  │ │ macros │ │ seeds  │ │ sources│  │
│  │  模型   │ │ 测试   │ │ 宏     │ │ 静态数据│ │ 数据源 │  │
│  └────┬───┘ └────┬───┘ └────┬───┘ └────┬───┘ └────┬───┘  │
│       │          │          │          │          │       │
│       └──────────┴──────────┴──────────┴──────────┘       │
│                          │ dbt compile/parse                 │
│                          ▼                                  │
│                   ┌─────────────┐                           │
│                   │  Jinja Template Engine                  │
│                   │  + 宏展开 + 配置解析                      │
│                   └──────┬──────┘                           │
│                          ▼                                  │
│                   ┌─────────────┐                           │
│                   │  Compiled SQL                           │
│                   │  (纯 ANSI SQL)                          │
│                   └──────┬──────┘                           │
│                          ▼                                  │
│              ┌───────────────────────┐                      │
│              │  DWH (Snowflake/BigQuery/Redshift/Doris)    │
│              │  实际执行 SQL 创建表/视图                     │
│              └───────────────────────┘                      │
└─────────────────────────────────────────────────────────────┘

1.3 项目目录结构

dbt_project/
├── dbt_project.yml          # 项目配置文件
├── profiles.yml             # 数据库连接配置(~/.dbt/)
├── packages.yml             # 依赖包
├── analyses/                # 分析查询(不物化)
│   └── ad_hoc_metrics.sql
├── macros/                  # 可复用 SQL 宏
│   ├── custom_tests.sql
│   └── utils.sql
├── models/                  # 核心数据模型
│   ├── staging/             # 原始数据清洗层 (stg_)
│   │   ├── stg_orders.sql
│   │   └── staging.yml      # 模型配置 + 测试
│   ├── marts/               # 业务主题层
│   │   ├── core/            # 核心模型
│   │   │   ├── dim_users.sql
│   │   │   └── fct_orders.sql
│   │   └── marketing/       # 营销域
│   │       └── mart_campaign_performance.sql
│   └── intermediate/        # 中间计算模型 (int_)
│       └── int_order_payments.sql
├── seeds/                   # 静态 CSV 数据
│   └── country_codes.csv
├── snapshots/               # 缓慢变化维快照
│   └── snapshot_users.sql
├── tests/                   # 自定义测试
│   └── assert_positive_revenue.sql
├── docs/                    # 额外文档
│   └── overview.md
└── snapshots/               # SCD Type 2 快照
    └── dim_users_snapshot.sql

2. Models 模型开发

2.1 Materialization 物化策略

物化类型特点适用场景
view不存储数据,实时查询轻量转换、维度表
table全量重建小表、每日全量
incremental增量追加/更新大表、流式数据
ephemeral内联 CTE中间计算、无需复用
materialized_view数据库物化视图数据库原生支持
-- models/staging/stg_orders.sql
{{ config(
    materialized='table',
    unique_key='order_id',
    partition_by={
      "field": "created_at",
      "data_type": "timestamp",
      "granularity": "day"
    }
) }}

WITH source AS (
    SELECT * FROM {{ source('raw', 'orders') }}
),
renamed AS (
    SELECT
        order_id,
        user_id,
        CAST(amount AS DECIMAL(18,2)) AS amount,
        status,
        created_at,
        updated_at,
        -- 派生字段
        CASE 
            WHEN amount >= 1000 THEN 'high'
            WHEN amount >= 100 THEN 'medium'
            ELSE 'low'
        END AS value_tier
    FROM source
    WHERE created_at >= '2020-01-01'
)

SELECT * FROM renamed

2.2 Incremental 增量模型

-- models/marts/fct_orders.sql
{{ config(
    materialized='incremental',
    unique_key='order_id',
    incremental_strategy='merge',  -- merge / insert_overwrite / append
    on_schema_change='sync_all_columns'
) }}

SELECT
    order_id,
    user_id,
    amount,
    status,
    created_at
FROM {{ ref('stg_orders') }}

{% if is_incremental() %}
    -- 增量条件:只处理新数据
    WHERE created_at > (SELECT MAX(created_at) FROM {{ this }})
{% endif %}

-- 运行方式:
-- dbt run --full-refresh    → 全量重建(DELETE + 全量 INSERT)
-- dbt run                   → 增量追加(只处理新增数据)

增量策略对比

策略数据库行为适用
mergeSnowflake/BigQuery/DatabricksMERGE INTO有 unique_key,支持更新
insert_overwriteBigQuery/Spark分区覆盖分区表,按分区全量替换
append通用INSERT INTO只有追加无更新
delete+insertDremio/StarRocks先删后插有限支持 merge 的引擎

2.3 Ref 与 Source

-- {{ source() }} 引用上游原始数据
{{ source('raw_schema', 'orders') }}
-- 在 models/sources.yml 中定义:
-- sources:
--   - name: raw_schema
--     tables:
--       - name: orders
--         loaded_at_field: created_at

-- {{ ref() }} 引用 dbt 模型
{{ ref('stg_orders') }}
-- dbt 自动解析依赖关系,构建 DAG
# models/sources.yml
version: 2

sources:
  - name: raw_ecommerce
    database: raw_db
    schema: ecommerce
    tables:
      - name: orders
        description: "原始订单数据"
        columns:
          - name: order_id
            description: "订单唯一标识"
            tests:
              - not_null
              - unique
        freshness:
          warn_after: {count: 12, period: hour}
          error_after: {count: 24, period: hour}
        loaded_at_field: created_at

      - name: users
        description: "原始用户数据"

3. Tests 数据质量测试

3.1 内置测试

测试类型说明示例
unique列值唯一tests: - unique
not_null列值非空tests: - not_null
accepted_values枚举值检查accepted_values: [1, 2, 3]
relationships外键关联to: ref('users'), field: user_id
# models/staging/staging.yml
version: 2

models:
  - name: stg_orders
    description: "清洗后的订单数据"
    columns:
      - name: order_id
        description: "主键"
        tests:
          - unique
          - not_null

      - name: user_id
        description: "用户ID,关联 dim_users"
        tests:
          - not_null
          - relationships:
              to: ref('stg_users')
              field: user_id

      - name: status
        description: "订单状态"
        tests:
          - accepted_values:
              values: ['pending', 'paid', 'shipped', 'delivered', 'cancelled']

      - name: amount
        description: "订单金额"
        tests:
          - not_null
          - dbt_utils.greater_than_or_equal_to:
              value: 0

3.2 自定义测试(Singular Test)

-- tests/assert_positive_revenue.sql
-- 单文件测试:检查总收入为正
SELECT
    SUM(amount) AS total_revenue
FROM {{ ref('fct_orders') }}
HAVING SUM(amount) < 0

-- 如果查询返回结果,则测试失败

3.3 通用测试(Generic Test)

-- macros/test_greater_than_or_equal_to.sql
{% test greater_than_or_equal_to(model, column_name, value) %}

SELECT
    {{ column_name }}
FROM {{ model }}
WHERE {{ column_name }} < {{ value }}

{% endtest %}
# 使用自定义测试
models:
  - name: stg_orders
    columns:
      - name: amount
        tests:
          - greater_than_or_equal_to:
              value: 0

4. Macros 与 Jinja 模板

4.1 Jinja 基础语法

-- 变量
{% set start_date = '2024-01-01' %}
WHERE dt >= '{{ start_date }}'

-- 逻辑控制
{% if target.name == 'prod' %}
    {{ config(materialized='table') }}
{% else %}
    {{ config(materialized='view') }}
{% endif %}

-- 循环
{% set payment_methods = ['credit_card', 'coupon', 'bank_transfer', 'gift_card'] %}

SELECT
    order_id,
    {% for payment_method in payment_methods %}
    SUM(CASE WHEN payment_method = '{{ payment_method }}' THEN amount ELSE 0 END)
        AS {{ payment_method }}_amount
    {%- if not loop.last -%},{% endif %}
    {% endfor %}
FROM {{ ref('stg_payments') }}
GROUP BY order_id

-- 生成:
-- SELECT order_id,
--   SUM(CASE WHEN payment_method = 'credit_card' THEN amount ELSE 0 END) AS credit_card_amount,
--   SUM(CASE WHEN payment_method = 'coupon' THEN amount ELSE 0 END) AS coupon_amount,
--   ...

4.2 实用宏

-- macros/generate_surrogate_key.sql
{% macro generate_surrogate_key(field_list) %}
    -- 生成代理键(类似 dbt_utils 的 generate_surrogate_key)
    CAST(
        MD5(
            CONCAT(
                {% for field in field_list %}
                    COALESCE(CAST({{ field }} AS STRING), '_dbt_null_')
                    {%- if not loop.last -%}, '||', {% endif %}
                {% endfor %}
            )
        ) AS STRING
    )
{% endmacro %}

-- 使用
SELECT
    {{ generate_surrogate_key(['user_id', 'order_date']) }} AS sk_order,
    *
FROM {{ ref('stg_orders') }}
-- macros/date_spine.sql
{% macro date_spine(datepart, start_date, end_date) %}
    WITH date_series AS (
        SELECT {{ dbt.dateadd(datepart, "n", "'" ~ start_date ~ "'") }} AS date_{{ datepart }}
        FROM UNNEST(GENERATE_ARRAY(0, 
            DATE_DIFF('{{ end_date }}', '{{ start_date }}', {{ datepart }})
        )) AS n
    )
    SELECT * FROM date_series
{% endmacro %}

5. Docs 文档生成

5.1 文档即代码

# models/docs.md
{% docs orders_status %}

订单状态说明:

| 状态 | 说明 |
|------|------|
| pending | 待支付 |
| paid | 已支付 |
| shipped | 已发货 |
| delivered | 已送达 |
| cancelled | 已取消 |

{% enddocs %}

# models/staging/staging.yml
version: 2
models:
  - name: stg_orders
    description: "订单明细表,来源 raw.orders"
    columns:
      - name: status
        description: '{{ doc("orders_status") }}'

5.2 自动生成文档站点

# 生成文档
# dbt docs generate

# 启动文档服务器
# dbt docs serve --port 8080

# 产出:
# - index.html: 搜索和导航
# - Lineage DAG: 表血缘关系图
# - Table details: 列定义、测试、描述

6. CI/CD 与 DataOps

6.1 dbt + GitHub Actions CI

# .github/workflows/dbt-ci.yml
name: dbt CI

on:
  pull_request:
    branches: [main]

jobs:
  dbt-ci:
    runs-on: ubuntu-latest
    steps:
      - uses: actions/checkout@v4
      
      - uses: actions/setup-python@v5
        with:
          python-version: '3.11'
      
      - name: Install dbt
        run: |
          pip install dbt-snowflake  # 或 dbt-bigquery, dbt-doris 等
          dbt deps
      
      - name: Compile
        run: dbt compile --target ci
      
      - name: Lint
        run: sqlfluff lint models/  # SQL 风格检查
      
      - name: Run tests on modified models
        run: |
          dbt run --select state:modified+ --target ci
          dbt test --select state:modified+ --target ci
        env:
          DBT_PROFILES_DIR: ./
          SNOWFLAKE_ACCOUNT: ${{ secrets.SNOWFLAKE_ACCOUNT }}
          SNOWFLAKE_USER: ${{ secrets.SNOWFLAKE_USER }}
          SNOWFLAKE_PASSWORD: ${{ secrets.SNOWFLAKE_PASSWORD }}

      - name: Generate docs
        run: dbt docs generate
      
      - name: Upload docs artifact
        uses: actions/upload-artifact@v4
        with:
          name: dbt-docs
          path: target/

6.2 Slim CI(只跑变更)

# 利用 dbt 的状态比较,只跑变更的模型
# dbt run --select state:modified+
# state:modified  = 相对于上一次 manifest 有变更的模型
# +               = 加上下游依赖

# 完整 Slim CI 流程
# 1. 下载 main 分支的 manifest.json
# 2. dbt compile --state path/to/main/manifest
# 3. dbt run --select state:modified+ --defer --state path/to/main/manifest
# 4. dbt test --select state:modified+

6.3 dbt Cloud vs dbt Core

特性dbt Core(开源)dbt Cloud(商业)
价格免费按席位收费
部署自管服务器/CI托管 SaaS
IDEVS Code + 插件内置 Web IDE
调度Airflow/Cron/自行内置 Scheduler
文档托管自行部署自动生成托管
CI/CDGitHub Actions 自建原生 Slim CI
监控告警自建内置
SSO/权限企业级 RBAC

7. 生产最佳实践

7.1 项目组织结构

-- dbt_project.yml 关键配置
name: 'ecommerce_analytics'
version: '1.0.0'
config-version: 2

profile: 'ecommerce'

model-paths: ["models"]
analysis-paths: ["analyses"]
test-paths: ["tests"]
seed-paths: ["seeds"]
macro-paths: ["macros"]
snapshot-paths: ["snapshots"]

target-path: "target"
clean-targets:
  - "target"
  - "dbt_packages"

models:
  ecommerce_analytics:
    staging:
      +materialized: view
      +schema: staging
    intermediate:
      +materialized: ephemeral
      +schema: intermediate
    marts:
      +materialized: table
      +schema: marts
      core:
        +tags: ["core", "daily"]
      marketing:
        +tags: ["marketing", "daily"]

7.2 环境隔离

# ~/.dbt/profiles.yml
ecommerce:
  target: dev
  outputs:
    dev:
      type: snowflake
      account: xyz.eu-west-1
      user: dbt_dev
      password: "{{ env_var('DBT_DEV_PASSWORD') }}"
      role: analyst
      database: analytics_dev
      warehouse: dev_wh
      schema: "{{ env_var('DBT_USER') }}"
      threads: 4

    prod:
      type: snowflake
      account: xyz.eu-west-1
      user: dbt_prod
      password: "{{ env_var('DBT_PROD_PASSWORD') }}"
      role: dbt_prod
      database: analytics
      warehouse: prod_wh
      schema: marts
      threads: 8

7.3 性能优化

-- 1. 分区与聚簇
{{ config(
    materialized='incremental',
    partition_by={
        "field": "created_at",
        "data_type": "timestamp",
        "granularity": "day"
    },
    cluster_by=["user_id", "status"]
) }}

-- 2. 查询标签追踪
{{ config(
    query_tag = 'dbt_model_' ~ this.name
) }}

-- 3. 预钩子和后钩子
{{ config(
    pre_hook="ALTER SESSION SET QUERY_TAG = 'dbt_{{this.name}}'",
    post_hook="GRANT SELECT ON {{ this }} TO ROLE analyst"
) }}

总结

实践推荐方案
模型组织staging → intermediate → marts
物化策略staging(view) + marts(table/incremental)
数据质量列级测试 + Singular 测试 + dbt-expectations
复用逻辑Macros + packages (dbt-utils, dbt-labs codegen)
代码复用dbt packages.yml 引用公共包
CI/CDGitHub Actions + state:modified Slim CI
文档docs blocks + dbt docs generate/serve
环境profiles.yml + target 切换

8. dbt Core vs dbt Cloud — 深度选型分析

8.1 功能全面对比

维度dbt Core(开源)dbt Cloud(SaaS)
价格免费,Apache 2.0Developer: $0/席;Team: $100/席/月;Enterprise: 按询价
部署方式自托管(本地/ECS/K8s/CI Runner)托管 SaaS,免运维
代码编辑器VS Code + dbt-power-user 插件 / Vim / IDE 任选内置 Web IDE,支持 AI Copilot、代码补全
调度执行Airflow / Dagster / Prefect / Cron 自建原生 Scheduler,支持 Cron 表达式
文档托管dbt docs generate + S3/Nginx 自托管自动生成,每次运行自动更新
CI/CDGitHub Actions / GitLab CI 自建原生 Slim CI,自动检测修改模型
监控告警依赖外部系统内置 Job 运行日志、失败告警、Slack 集成
SSO / RBACEnterprise 支持 SAML/Okta/SSO,细粒度权限
Semantic LayerEnterprise 支持统一语义层

8.2 自托管 vs SaaS 选型决策树

场景推荐方案理由
初创团队 < 5 人,快速启动dbt Cloud Developer零运维,内置调度,快速验证
中型团队 5-20 人,已有 Airflowdbt Core + Airflow已有调度体系,避免重复投入
金融/医疗等行业,数据不出境dbt Core 私有化合规要求,需完整控制运行环境
大型企業 50+ 分析师,需统一治理dbt Cloud EnterpriseSSO、RBAC、审计日志、SLA 保障
本地开发/POC 验证dbt Core + DuckDB零外部依赖,笔记本即可运行

8.3 CI/CD 集成差异

dbt Core 需要自行搭建完整的 CI 流水线,包括安装依赖、编译检查、状态对比、增量运行等。dbt Cloud 则提供开箱即用的 Slim CI,只需连接 Git 仓库即可自动触发。


9. dbt 高级特性

9.1 Jinja 宏包管理(packages.yml)

dbt 支持通过 packages.yml 引入社区宏包,实现代码复用。

# packages.yml
packages:
  - package: dbt-labs/dbt_utils
    version: 1.1.1

  - package: calogica/dbt_expectations
    version: 0.10.3

  - package: dbt-labs/codegen
    version: 0.12.1

  - git: "https://github.com/my-org/dbt-common-macros.git"
    revision: v1.2.0
# 安装依赖包
$ dbt deps

9.2 种子 Seeds(静态数据管理)

Seeds 用于管理不经常变更的静态数据,如国家代码、汇率映射、产品分类等。

# seeds/country_codes.csv
country_code,country_name,region,currency
US,United States,North America,USD
CN,China,Asia,CNY
DE,Germany,Europe,EUR
JP,Japan,Asia,JPY
GB,United Kingdom,Europe,GBP
-- 在模型中引用 seed
SELECT
    o.*,
    c.country_name,
    c.currency
FROM {{ ref('stg_orders') }} o
LEFT JOIN {{ ref('country_codes') }} c
    ON o.country_code = c.country_code

9.3 快照 Snapshots(SCD Type 2)

快照自动追踪记录历史变更,实现缓慢变化维(Slowly Changing Dimension Type 2)。

-- snapshots/snapshot_users.sql
{% snapshot snapshot_users %}

{{
    config(
      target_database='analytics',
      target_schema='snapshots',
      unique_key='user_id',
      strategy='timestamp',
      updated_at='updated_at',
    )
}}

SELECT * FROM {{ source('raw_ecommerce', 'users') }}

{% endsnapshot %}

快照输出表自动包含四个元字段:

字段说明
dbt_scd_id变更记录的代理键
dbt_updated_at快照执行时间
dbt_valid_from记录生效起始时间
dbt_valid_to记录失效时间(当前记录为 NULL)

9.4 增量模型策略详解

-- merge 策略(Snowflake/BigQuery/Databricks)
{{ config(
    materialized='incremental',
    unique_key='order_id',
    incremental_strategy='merge',
    merge_update_columns=['status', 'amount', 'updated_at']
) }}

SELECT * FROM {{ source('raw_ecommerce', 'orders') }}
{% if is_incremental() %}
    WHERE updated_at > (SELECT MAX(updated_at) FROM {{ this }})
{% endif %}
-- insert_overwrite 策略(适合 BigQuery 分区表)
{{ config(
    materialized='incremental',
    partition_by={"field": "created_at", "data_type": "date"},
    incremental_strategy='insert_overwrite'
) }}

-- 每次运行覆盖最近 7 天分区的全量数据
SELECT * FROM {{ source('raw_ecommerce', 'orders') }}
WHERE created_at >= DATE_SUB(CURRENT_DATE(), INTERVAL 7 DAY)
-- delete+insert 策略(Dremio/StarRocks)
{{ config(
    materialized='incremental',
    unique_key='order_id',
    incremental_strategy='delete+insert'
) }}

10. dbt 测试体系

10.1 内置测试清单

测试类型作用适用层级
unique列值全局唯一主键列
not_null列值非空业务关键列
relationships外键引用完整性关联列
accepted_values枚举值白名单状态/类型列
row_count行数阈值检查(dbt-core v1.8+)模型级
# models/marts/marts.yml
version: 2

models:
  - name: fct_orders
    columns:
      - name: order_id
        tests:
          - unique
          - not_null

      - name: status
        tests:
          - accepted_values:
              values: ['pending', 'paid', 'shipped', 'delivered', 'cancelled']

      - name: user_id
        tests:
          - relationships:
              to: ref('dim_users')
              field: user_id

      - name: amount
        tests:
          - dbt_utils.greater_than_or_equal_to:
              value: 0

10.2 自定义通用测试(Generic Test)

-- macros/test_column_is_date.sql
{% test column_is_date(model, column_name) %}

SELECT
    {{ column_name }}
FROM {{ model }}
WHERE SAFE.PARSE_DATE('%Y-%m-%d', CAST({{ column_name }} AS STRING)) IS NULL
  AND {{ column_name }} IS NOT NULL

{% endtest %}

10.3 自定义单例测试(Singular Test)

-- tests/assert_payment_method_total.sql
SELECT
    payment_method,
    SUM(amount) AS total_amount
FROM {{ ref('fct_orders') }}
GROUP BY payment_method
HAVING total_amount < 0

10.4 Great Expectations 集成

# packages.yml
packages:
  - package: calogica/dbt_expectations
    version: 0.10.3
# 使用 dbt_expectations 测试
models:
  - name: fct_orders
    columns:
      - name: amount
        tests:
          - dbt_expectations.expect_column_values_to_be_between:
              min_value: 0
              max_value: 100000

      - name: created_at
        tests:
          - dbt_expectations.expect_row_values_to_have_recent_data:
              datepart: day
              interval: 1

11. dbt 文档与血缘

11.1 自动文档生成

dbt 自动从模型定义、列注释、测试配置中提取元数据,生成可搜索的文档站点。

# 生成静态文档到 target/ 目录
$ dbt docs generate

# 本地预览文档
$ dbt docs serve --port 8080

11.2 DAG 可视化与导航

文档站点内置交互式 DAG,支持:

  • 上游/下游血缘展开
  • 列级血缘追踪(dbt Cloud + dbt Explorer)
  • 快速跳转到模型 SQL 源码

11.3 暴露 Exposures

Exposures 声明下游消费方,如 BI 报表、ML 模型、反向 ETL。

# models/exposures.yml
version: 2

exposures:
  - name: executive_dashboard
    type: dashboard
    maturity: high
    url: https://looker.company.com/dashboards/123
    description: "高管日报仪表板,核心经营指标汇总"
    depends_on:
      - ref('fct_orders')
      - ref('fct_revenue_daily')
    owner:
      name: "数据产品团队"
      email: "data-product@company.com"

11.4 Metrics 语义层(v1 与 v2)

# metrics/revenue_metrics.yml (dbt Metrics v1 / MetricFlow v2)
version: 2

metrics:
  - name: total_revenue
    label: "总营收"
    model: ref('fct_orders')
    description: "所有已完成订单的金额总和"
    calculation_method: sum
    expression: amount
    timestamp: created_at
    time_grains: [day, week, month, quarter, year]
    dimensions:
      - status
      - country_code
    filters:
      - field: status
        operator: '!='
        value: 'cancelled'

12. dbt 与 Airflow 集成

12.1 BashOperator 直接调用 dbt-core

# airflow/dags/dbt_pipeline.py
from airflow import DAG
from airflow.operators.bash import BashOperator
from airflow.utils.dates import days_ago

with DAG(
    dag_id='dbt_daily_pipeline',
    start_date=days_ago(1),
    schedule_interval='0 3 * * *',
    catchup=False
) as dag:

    dbt_deps = BashOperator(
        task_id='dbt_deps',
        bash_command='cd /opt/dbt && dbt deps --profiles-dir .'
    )

    dbt_run = BashOperator(
        task_id='dbt_run',
        bash_command='cd /opt/dbt && dbt run --target prod --profiles-dir .'
    )

    dbt_test = BashOperator(
        task_id='dbt_test',
        bash_command='cd /opt/dbt && dbt test --target prod --profiles-dir .'
    )

    dbt_deps >> dbt_run >> dbt_test

12.2 dbt Cloud Job API 触发

# 通过 dbt Cloud API 触发 Job
import requests

account_id = 12345
job_id = 67890
token = "dbtu_xxxxxxxx"

response = requests.post(
    f"https://cloud.getdbt.com/api/v2/accounts/{account_id}/jobs/{job_id}/run/",
    headers={"Authorization": f"Token {token}"},
    json={"cause": "Triggered by Airflow"}
)
# Airflow 中使用 SimpleHttpOperator
from airflow.providers.http.operators.http import SimpleHttpOperator

trigger_dbt_cloud_job = SimpleHttpOperator(
    task_id='trigger_dbt_cloud_job',
    http_conn_id='dbt_cloud_api',
    endpoint='api/v2/accounts/{{ var.value.dbt_account_id }}/jobs/{{ var.value.dbt_job_id }}/run/',
    method='POST',
    headers={"Content-Type": "application/json"},
    data='{"cause": "Airflow triggered"}'
)

12.3 Astronomer Cosmos 框架

Cosmos 自动将 dbt 项目映射为 Airflow DAG 或 TaskGroup,无需手动维护任务依赖。

# airflow/dags/cosmos_dag.py
from cosmos import DbtDag, ProjectConfig, ProfileConfig
from cosmos.profiles import SnowflakeUserPasswordProfileMapping

profile_config = ProfileConfig(
    profile_name="ecommerce",
    target_name="prod",
    profile_mapping=SnowflakeUserPasswordProfileMapping(
        conn_id="snowflake_conn",
        profile_args={"database": "analytics", "schema": "marts"}
    )
)

dbt_snowflake_dag = DbtDag(
    project_config=ProjectConfig("/opt/dbt/ecommerce_analytics"),
    profile_config=profile_config,
    dag_id="cosmos_dbt_pipeline",
    schedule_interval="0 3 * * *",
    start_date=days_ago(1),
    catchup=False,
    default_args={"retries": 2}
)

13. dbt 项目结构最佳实践

13.1 目录规范

dbt_project/
├── models/
│   ├── staging/              # 原始数据清洗,轻量转换
│   │   ├── _sources.yml      # source 定义
│   │   ├── _staging.yml      # staging 层模型配置 + 测试
│   │   ├── stg_orders.sql
│   │   └── stg_users.sql
│   ├── intermediate/         # 中间计算模型,面向分析主题
│   │   ├── int_orders_enriched.sql
│   │   └── int_user_sessions.sql
│   └── marts/                # 业务主题模型,下游直接使用
│       ├── core/
│       │   ├── dim_users.sql
│       │   ├── fct_orders.sql
│       │   └── _core.yml
│       └── marketing/
│           ├── mart_campaign_performance.sql
│           └── _marketing.yml

13.2 命名约定

前缀含义物化策略
stg_原始数据清洗模型view
int_中间计算模型ephemeral / view
dim_维度表table / incremental
fct_事实表table / incremental
agg_聚合表table
rpt_报表输出table

13.3 模型分层职责

层级职责变更频率下游依赖
Sources声明上游系统表结构随源系统变更仅 staging
Staging字段重命名、类型转换、简单过滤全层
Intermediate多表 join、业务逻辑计算marts
Marts主题域建模,供 BI/ML 直接使用外部系统

13.4 环境隔离

# ~/.dbt/profiles.yml
my_project:
  target: dev
  outputs:
    dev:
      type: bigquery
      method: oauth
      project: my-gcp-project
      dataset: "dbt_{{ env_var('USER') }}"   # 每人独立 schema
      threads: 4

    staging:
      type: bigquery
      method: service-account
      project: my-gcp-project
      dataset: staging
      threads: 8

    prod:
      type: bigquery
      method: service-account
      project: my-gcp-project
      dataset: analytics
      threads: 16

14. dbt 性能优化

14.1 物化策略选择矩阵

数据量变更模式查询频率推荐策略
< 10 MB全量替换低频view
10 MB - 1 GB全量替换中频table
> 1 GB仅追加高频incremental + append
> 1 GB增删改高频incremental + merge
> 100 GB仅追加高频incremental + insert_overwrite(分区表)

14.2 Incremental 模型适用场景

增量模型适合以下场景:

  • 事件流数据,按时间分片持续追加
  • 日志类表,数据只增不删
  • 大表全量重建成本过高

不适合的场景:

  • 数据需要频繁回溯修正
  • 源数据会物理删除历史记录
  • 每日数据量变化极小(< 100 MB),全量重建更快

14.3 集群键与分区键配合

-- BigQuery 分区 + 聚簇优化
{{ config(
    materialized='incremental',
    partition_by={
        "field": "created_at",
        "data_type": "timestamp",
        "granularity": "day"
    },
    cluster_by=["user_id", "status"]
) }}

SELECT
    order_id,
    user_id,
    amount,
    status,
    created_at
FROM {{ source('raw_ecommerce', 'orders') }}
WHERE created_at >= TIMESTAMP_SUB(CURRENT_TIMESTAMP(), INTERVAL 7 DAY)

14.4 查询下推与并行度调优

# dbt_project.yml
models:
  my_project:
    +pre-hook:
      - "SET query_tag = 'dbt_{{ this.name }}'"
    +post-hook:
      - "GRANT SELECT ON {{ this }} TO ROLE analyst"
    staging:
      +snowflake_warehouse: "DEV_WH"
      +threads: 4
    marts:
      +snowflake_warehouse: "PROD_WH"
      +threads: 8
# 调整单命令并行度
$ dbt run --threads 16              # 默认 4,可提升至数据仓库上限
$ dbt run --select marts.fct_orders --threads 8

15. dbt 与 Data Lake 集成

15.1 dbt-spark 适配器(Spark / Iceberg / Delta Lake)

# profiles.yml
spark_datalake:
  target: dev
  outputs:
    dev:
      type: spark
      method: thrift
      host: spark-thrift-server.hadoop
      port: 10000
      schema: analytics
      threads: 4
-- 使用 Delta Lake 增量模型
{{ config(
    materialized='incremental',
    file_format='delta',
    incremental_strategy='merge',
    unique_key='order_id'
) }}

SELECT * FROM {{ source('bronze', 'orders') }}
{% if is_incremental() %}
    WHERE modified_at > (SELECT MAX(modified_at) FROM {{ this }})
{% endif %}

15.2 dbt-duckdb 本地分析

DuckDB 是无外部依赖的嵌入式分析数据库,适合本地开发与测试。

# profiles.yml
local_analytics:
  target: dev
  outputs:
    dev:
      type: duckdb
      path: '/tmp/analytics.duckdb'
      threads: 4
# 本地极速运行,无需网络连接
$ dbt run --target dev --select staging
$ dbt test --target dev

15.3 dbt-trino 联邦查询

Trino 支持跨多个数据源的联邦查询,dbt-trino 适配器可将 dbt 模型物化到 Trino 连接的任意后端。

# profiles.yml
federation:
  target: prod
  outputs:
    prod:
      type: trino
      host: trino.company.com
      port: 8080
      user: dbt
      catalog: iceberg
      schema: marts
      threads: 8

16. 常见问题(FAQ)

Q1: dbt 中的 ephemeral 模型和普通 view 有什么区别?
ephemeral 不会物化为数据库对象,而是以内联 CTE 的形式被引用它的模型编译进去。适合中间计算逻辑,不需要单独被查询。view 则会在数据库中创建真正的视图对象,可被任意 SQL 查询。

Q2: 增量模型如何处理源数据的历史修正?
增量模型默认只处理新数据。如果源数据可能回溯修正历史记录,需要使用 merge 策略配合合适的 unique_key,确保历史记录被覆盖更新。对于大规模回溯,建议执行 dbt run --full-refresh 进行全量重建。

Q3: dbt Core 如何实现数据质量告警?
dbt Core 本身不提供告警系统。常见的方案是在 CI/CD 中捕获 dbt test 的退出码(非零即失败),通过 GitHub Actions / Slack Webhook / PagerDuty 发送告警。dbt Cloud 则内置失败告警和 Slack 集成。

Q4: snapshots 和 incremental 模型都可以追踪变更,该如何选择?
snapshots 用于追踪源数据的历史变更(谁改了什么、何时改的),输出 SCD Type 2 表。incremental 模型用于增量构建派生表,不关心上游历史。两者目的不同,通常组合使用:先 snapshot 保留历史,再基于 snapshot 做增量 marts。

Q5: dbt 是否适合数据量极小的项目(< 1 GB)?
即使数据量很小,dbt 提供的版本控制、测试框架、文档生成、环境隔离等工程化能力仍然具有价值。对于极轻量场景,可以配合 dbt-duckdb 实现零基础设施的本地分析,后续随数据增长无缝迁移到 Snowflake/BigQuery。


总结

dbt 将软件工程的最佳实践引入数据转换领域,彻底改变了传统 SQL 脚本的管理方式。从 dbt Core 的开源生态到 dbt Cloud 的托管服务,从基础的模型物化到高级的语义层(Metrics)和快照(Snapshots),dbt 为不同规模的数据团队提供了完整的 DataOps 解决方案。

关键成功要素包括:

  • 建立清晰的模型分层(staging / intermediate / marts)
  • 为每个模型配置适当的物化策略与测试
  • 利用 Jinja 宏和标准包实现代码复用
  • 将 dbt 集成到 CI/CD 和 Airflow 调度中
  • 借助自动文档和 DAG 血缘提升数据可发现性

无论团队使用 Snowflake、BigQuery、Databricks、Spark 还是 DuckDB,dbt 都能提供统一的抽象层,让数据分析师用纯 SQL 构建可靠、可维护、可扩展的数据管道。

继续阅读

探索更多技术文章

浏览归档,发现更多关于系统设计、工具链和工程实践的内容。

全部文章 返回首页

「data-engineering」更多文章

  1. 数据工程深度指南:Modern Data Stack 全栈实践
  2. 数据平台工程:Data Mesh、FinOps 与 DataOps 生产实践
  3. Kafka Connect CDC 实战:Debezium 数据同步与变更捕获