Snowflake 上的 dbt
源端 Medallion 管道是如何搭建的:sources、staging 视图、增量 MERGE 模型、snapshot 与测试。
本文说明 dbt(data build tool)在 NYC Taxi Snowflake 迁移实验课中的用法,它做什么、每个部分为什么存在,以及应该如何理解它。
dbt 做什么(以及不做什么)
dbt 转换已经在数据库中的数据。它不会从外部加载数据、不搬运文件,也不管理基础设施。它的职责是把原始表变成干净、经过测试、可直接用于分析的表,通过运行你编写的 SQL 来实现。
可以把它看作 SQL 的构建系统。models/ 中的每个 .sql 文件都是一个模型,会在 Snowflake 中变成一张表或一个视图。dbt 负责处理 CREATE OR REPLACE 这类样板代码、解析模型之间的依赖关系,并运行你的测试。
项目结构
dbt/nyc_taxi_dbt/
├── dbt_project.yml # Project config: name, folder layout, materialization defaults
├── profiles.yml.example # Connection config template (copy to ~/.dbt/profiles.yml)
├── packages.yml # Third-party dbt packages
│
├── macros/
│ ├── generate_schema_name.sql # Overrides dbt's default schema naming logic
│ └── generate_surrogate_key.sql # Wrapper for consistent surrogate key generation
│
└── models/
├── sources.yml # Declares RAW.TRIPS_RAW as an external source
│
├── staging/ # Layer 1: clean and rename raw columns
│ ├── schema.yml # Column-level tests for staging models
│ ├── stg_trips.sql
│ └── stg_taxi_zones.sql
│
├── intermediate/ # Layer 2: joins and enrichment (no physical table)
│ └── int_trips_enriched.sql
│
└── analytics/ # Layer 3: final tables consumed by dashboards
├── schema.yml
├── fact_trips.sql
├── agg_hourly_zone_trips.sql
├── dim_date.sql
├── dim_payment_type.sql
├── dim_taxi_zones.sql
└── dim_vendor.sql三个层次(Medallion 架构)
第 1 层:Staging(models/staging/)
目的: 接收原样到达的原始数据,并使其可用。
这些模型以视图的形式落在 STAGING schema 中(没有存储成本,它们在查询时运行)。每个 staging 模型只做一件事:
| 模型 | 来源 | 它做什么 |
|---|---|---|
stg_trips | RAW.TRIPS_RAW | 将列名改为 snake_case,添加 duration_minutes,把 VARIANT 类型的 TRIP_METADATA 列展开为带类型的列 |
stg_taxi_zones | ANALYTICS.DIM_TAXI_ZONES | 轻量清洗,添加 COALESCE 保护,并提供一个 dbt 血缘节点 |
这里最重要的工作是展开 TRIP_METADATA VARIANT 列。Snowflake 的冒号路径语法可以提取嵌套的 JSON 字段:
-- Snowflake: colon-path notation
TRIP_METADATA:driver.rating::FLOAT AS driver_rating,
TRIP_METADATA:app.surge_multiplier::FLOAT AS surge_multiplier这正是迁移挑战之一,ClickHouse 改用 JSONExtractFloat(TRIP_METADATA, 'driver', 'rating')。
第 2 层:Intermediate(models/intermediate/)
目的: 把所有 join 集中在一处,避免重复编写。
int_trips_enriched 将 stg_trips 与每个维度(zones、payment types、vendors、dates)关联,为每次行程生成一行宽的、完全反范式化的记录。它被声明为 ephemeral,也就是说 dbt 会把它的 SQL 内联到引用它的模型中,在 Snowflake 中不会创建任何物理表或视图。
-- dbt_project.yml
intermediate:
+materialized: ephemeral # compiled inline, no CREATE TABLE当中间结果只被一个下游模型需要,并且你不想为存储或查询编译开销付费时,就使用 ephemeral。
第 3 层:Analytics(models/analytics/)
目的: 最终的、可直接用于看板的表。
这些落在 ANALYTICS schema 中。共有两类:
静态维度表,体量小,每次 dbt run 时完全重载:
| 模型 | 行数 | 说明 |
|---|---|---|
dim_date | ~7,670 | 2009–2029 的日期骨架,含财季和美国联邦假日 |
dim_payment_type | 6 | 由 seed 数据直通而来 |
dim_vendor | 3 | 由 seed 数据直通而来 |
dim_taxi_zones | 265 | 经 stg_taxi_zones 直通而来 |
增量事实表/聚合表,体量大,每次运行用 MERGE 更新:
| 模型 | 行数 | 说明 |
|---|---|---|
fact_trips | 50M | 每次行程一行,完全反范式化 |
agg_hourly_zone_trips | ~9M | 按 zone 预聚合的小时级计数 |
物化方式(Materializations)
物化方式决定 dbt 为某个模型在 Snowflake 中创建什么。
| 物化方式 | 对应的 Snowflake 对象 | 何时使用 |
|---|---|---|
view | CREATE VIEW | 便宜;始终反映最新数据;用于 staging |
table | CREATE TABLE AS SELECT | 每次运行完全重建;用于小型维度表 |
incremental | 对已有表执行 MERGE INTO | 大表;只处理新增行 |
ephemeral | (无对象,内联为 CTE) | 只被一个下游模型共享的中间逻辑 |
两个增量模型演示了不同的增量策略:
fact_trips,处理自上次运行以来的新行程:
{% if is_incremental() %}
WHERE pickup_at > (SELECT MAX(pickup_at) FROM {{ this }})
{% endif %}agg_hourly_zone_trips,对滚动的 2 小时窗口重新聚合,以捕获迟到数据:
{% if is_incremental() %}
WHERE pickup_at >= DATEADD('hour', -2, CURRENT_TIMESTAMP())
{% endif %}在第一次运行时(表为空),is_incremental() 返回 false,会处理完整数据集。后续运行只处理新数据。如果 schema 发生变化,你需要从头重建,运行:
dbt run --full-refreshMERGE 策略(关键迁移挑战)
当 incremental_strategy = 'merge' 时,dbt 会生成一条 Snowflake 的 MERGE INTO 语句:
MERGE INTO ANALYTICS.FACT_TRIPS AS target
USING (SELECT ...) AS source
ON target.trip_id = source.trip_id
WHEN MATCHED THEN UPDATE SET ...
WHEN NOT MATCHED THEN INSERT ...;这是本实验课记录的最重要的迁移挑战之一。ClickHouse 没有 MERGE 语句。 在 ClickHouse 中的等价做法是使用 ReplacingMergeTree 表引擎并在查询中加上 FINAL,或者使用 CollapsingMergeTree 来获得显式的插入/删除语义。
Schema 命名:generate_schema_name 宏
dbt 的默认行为是把 profiles.yml 中的目标 schema 与 dbt_project.yml 中的自定义 schema 拼接起来:
target schema = STAGING + custom schema = ANALYTICS → STAGING_ANALYTICS (wrong)本项目用 macros/generate_schema_name.sql 中的自定义宏覆盖了该行为:
{% macro generate_schema_name(custom_schema_name, node) -%}
{%- if custom_schema_name is none -%}
{{ target.schema | upper }} -- no custom schema → use target schema
{%- else -%}
{{ custom_schema_name | upper }} -- custom schema → use it directly
{%- endif -%}
{%- endmacro %}结果:带有 +schema: ANALYTICS 的模型会落在 ANALYTICS 中,而不是 STAGING_ANALYTICS。
只要你在一个 dbt 项目中有多个 schema,并且不希望目标 schema 名被前置拼接,就需要这个宏。
连接与凭据(profiles.yml)
dbt 通过定义在 ~/.dbt/profiles.yml 中的 profile 连接 Snowflake(该文件绝不提交到 git)。dbt_project.yml 中的 profile 名称必须一致:
# dbt_project.yml
profile: 'nyc_taxi'
# ~/.dbt/profiles.yml
nyc_taxi:
target: dev
outputs:
dev:
type: snowflake
account: "{{ env_var('SNOWFLAKE_ORG') }}-{{ env_var('SNOWFLAKE_ACCOUNT') }}"
role: DBT_ROLE
database: NYC_TAXI_DB
warehouse: TRANSFORM_WH
schema: STAGING # ← this is the "target schema" / default schema
threads: 4要点:
schema: STAGING是默认 schema。没有+schema:覆盖的模型都落在这里。role: DBT_ROLE是由 Terraform 创建的最小权限角色,只拥有 dbt 所需的权限。threads: 4控制 dbt 并行构建多少个模型。- 凭据来自环境变量,在运行 setup 之前从
.env加载。
测试
dbt 测试有两种形式:
Schema 测试(在 schema.yml 中声明)
- name: trip_id
tests:
- not_null
- unique
- name: total_amount_usd
tests:
- dbt_expectations.expect_column_values_to_be_between:
min_value: 0
max_value: 1000not_null 和 unique 是内置的。dbt_expectations 测试来自 packages.yml 中声明的 calogica/dbt_expectations 包。
自定义 SQL 测试(tests/)
-- tests/assert_revenue_positive.sql
-- A passing test returns 0 rows
SELECT trip_id, total_amount_usd
FROM {{ ref('fact_trips') }}
WHERE total_amount_usd < 0自定义测试就是普通的 SQL 查询。dbt 会运行它们,并在返回任何行时判定失败。
运行全部测试:
dbt test第三方包(packages.yml)
packages:
- package: dbt-labs/dbt_utils
version: [">=1.0.0", "<2.0.0"]
- package: calogica/dbt_expectations
version: [">=0.10.0", "<1.0.0"]首次使用前先安装:
dbt depsdbt_utils 提供了 dim_date.sql 中使用的 date_spine 生成器。dbt_expectations 提供了超出内置 not_null/unique 的范围与分布类测试。
依赖图
dbt 通过跟踪 {{ ref() }} 调用自动按正确顺序构建模型:
RAW.TRIPS_RAW (source — not managed by dbt)
└── stg_trips (view)
└── int_trips_enriched (ephemeral)
├── fact_trips (incremental table)
└── agg_hourly_zone_trips (incremental table)
ANALYTICS.DIM_TAXI_ZONES (seeded by SQL script)
└── stg_taxi_zones (view)
├── int_trips_enriched
└── dim_taxi_zones (table)
dbt_utils.date_spine
└── dim_date (table){{ ref('stg_trips') }} 是一个模型声明对另一个模型依赖的方式。{{ source('raw', 'TRIPS_RAW') }} 声明对外部表的依赖(在 sources.yml 中定义)。
常用命令
| 命令 | 作用 |
|---|---|
dbt deps | 安装 packages.yml 中的包 |
dbt run | 构建所有模型(尽可能走增量) |
dbt run --full-refresh | 从头重建所有增量模型 |
dbt run -s fact_trips | 只构建 fact_trips 及其依赖 |
dbt test | 运行所有 schema 测试与自定义测试 |
dbt build | dbt run + dbt test 一起执行 |
dbt compile | 生成 SQL 但不执行(便于调试) |
dbt docs generate && dbt docs serve | 构建并在浏览器中浏览血缘图 |
在本项目中,如果 fact_trips 为空(首次运行或拆除之后),setup.sh 会自动触发 dbt run --full-refresh。
dbt 在完整搭建流程中的位置
terraform apply → creates warehouses, database, schemas, roles
scripts/01_create_tables.sql → creates raw tables, seeds dimension data
scripts/02_seed_data.sql → loads 50M synthetic trip rows
dbt deps && dbt build → transforms raw data into analytics-ready tables
scripts/03_create_streams_tasks.sql → creates CDC stream and scheduled taskdbt 位于管道的中间。它必须等到原始表存在且有数据之后才能运行。setup.sh 脚本负责处理这一顺序。