Snowflake MigrationClickHouse Workshops

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_tripsRAW.TRIPS_RAW将列名改为 snake_case,添加 duration_minutes,把 VARIANT 类型的 TRIP_METADATA 列展开为带类型的列
stg_taxi_zonesANALYTICS.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,6702009–2029 的日期骨架,含财季和美国联邦假日
dim_payment_type6由 seed 数据直通而来
dim_vendor3由 seed 数据直通而来
dim_taxi_zones265经 stg_taxi_zones 直通而来

增量事实表/聚合表,体量大,每次运行用 MERGE 更新:

模型行数说明
fact_trips50M每次行程一行,完全反范式化
agg_hourly_zone_trips~9M按 zone 预聚合的小时级计数

物化方式(Materializations)

物化方式决定 dbt 为某个模型在 Snowflake 中创建什么。

物化方式对应的 Snowflake 对象何时使用
viewCREATE VIEW便宜;始终反映最新数据;用于 staging
tableCREATE 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-refresh

MERGE 策略(关键迁移挑战)

当 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: 1000

not_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 deps

dbt_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 builddbt 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 task

dbt 位于管道的中间。它必须等到原始表存在且有数据之后才能运行。setup.sh 脚本负责处理这一顺序。

本页内容

ZH