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 → 增量追加(只处理新增数据)
增量策略对比:
| 策略 | 数据库 | 行为 | 适用 |
|---|---|---|---|
| merge | Snowflake/BigQuery/Databricks | MERGE INTO | 有 unique_key,支持更新 |
| insert_overwrite | BigQuery/Spark | 分区覆盖 | 分区表,按分区全量替换 |
| append | 通用 | INSERT INTO | 只有追加无更新 |
| delete+insert | Dremio/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 |
| IDE | VS Code + 插件 | 内置 Web IDE |
| 调度 | Airflow/Cron/自行 | 内置 Scheduler |
| 文档托管 | 自行部署 | 自动生成托管 |
| CI/CD | GitHub 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/CD | GitHub 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.0 | Developer: $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/CD | GitHub Actions / GitLab CI 自建 | 原生 Slim CI,自动检测修改模型 |
| 监控告警 | 依赖外部系统 | 内置 Job 运行日志、失败告警、Slack 集成 |
| SSO / RBAC | 无 | Enterprise 支持 SAML/Okta/SSO,细粒度权限 |
| Semantic Layer | 无 | Enterprise 支持统一语义层 |
8.2 自托管 vs SaaS 选型决策树
| 场景 | 推荐方案 | 理由 |
|---|---|---|
| 初创团队 < 5 人,快速启动 | dbt Cloud Developer | 零运维,内置调度,快速验证 |
| 中型团队 5-20 人,已有 Airflow | dbt Core + Airflow | 已有调度体系,避免重复投入 |
| 金融/医疗等行业,数据不出境 | dbt Core 私有化 | 合规要求,需完整控制运行环境 |
| 大型企業 50+ 分析师,需统一治理 | dbt Cloud Enterprise | SSO、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 构建可靠、可维护、可扩展的数据管道。
继续阅读
探索更多技术文章
浏览归档,发现更多关于系统设计、工具链和工程实践的内容。