Snowflake MigrationClickHouse Workshops

ClickHouse 上的 dbt

配置 dbt-clickhouse:delete_insert 增量策略、ReplacingMergeTree 模型,以及可刷新的 materialized view。

本指南介绍你将在第 3 部分用到的 dbt-clickhouse 专属模式。请在完成练习册 1–4 之后、开始练习册 5(dbt 模型设计)之前阅读。

如果你是从 dbt-snowflake 过来的,大部分 dbt 概念是完全一致的,sources、refs、测试、宏,以及 staging/intermediate/analytics 的分层模式。变化的是 ClickHouse 专属的配置层:engine、order_by、增量策略,以及 FINAL 语义。


1. 物化类型

dbt-clickhouse 支持五种物化方式。选择依据是更新模式,而不是个人偏好。

物化方式物理对象何时使用
viewClickHouse 视图Staging 模型:清洗并转换源数据类型;无存储成本;每次查询时重新计算
ephemeral无对象(内联为 CTE)通过 JOIN 组合多个 staging 模型的中间模型;避免创建冗余的物理表
table先在一个临时关系中构建完整的替换版本,再通过 EXCHANGE TABLES(或在旧版本上用成对重命名)原子地换入;换入后删除旧表每次 dbt run 都完全替换的小型维度表;不需要局部更新。注意: 对大表来说完全重建不可行,任何超过几千行的表都应使用 incremental。
incremental首次运行时 CREATE TABLE;后续运行采用选择性 UPDATE 模式事实表和预聚合表,每次运行只需处理新增/变更的行
materialized_viewClickHouse Materialized View自动刷新的聚合;与 dbt 的 incremental 不同。标准(基于触发器的)MV 在每次 INSERT 时触发一次,并且只能看到那一批数据,它无法计算全生命周期的聚合。而 REFRESHABLE MV 会按计划重新运行整个查询,因此可以做到。

与 Snowflake 的关键差异: dbt-snowflake 在内部处理存储细节。在 dbt-clickhouse 中,table 和 incremental 模型需要显式的 +engine 配置,dbt 用它来生成 CREATE TABLE ... ENGINE = ... DDL。

视图没有 engine。 如果你不小心给 view 物化加上了 +engine,dbt-clickhouse 会忽略它。只有 table 和 incremental 物化才会创建需要 engine 的持久化存储。

可刷新的 materialized view。 dbt-clickhouse 的 materialized_view 物化接受一个 refreshable 配置块,包含 interval(以及可选的 randomize),它会在生成的 CREATE MATERIALIZED VIEW 语句中直接输出 REFRESH 子句。本实验课的 mv_live_trip_feed 模型没有设置 refreshable,这就是它构建出的 MV 没有刷新计划的原因。


2. 在 dbt 中表达 ClickHouse 配置

ClickHouse 专属设置以 dbt 模型配置的形式表达,可以写在 dbt_project.yml 中(作为项目级默认值),也可以写在某个模型的 config() 块中(作为模型级覆盖)。

在 dbt_project.yml 中

models:
  your_project:
    analytics:
      +schema: analytics
      +materialized: table
      +engine: "MergeTree()"          # default for all analytics tables

      fact_trips:
        +materialized: incremental
        +engine: "ReplacingMergeTree(updated_at)"   # overrides the default
        +incremental_strategy: delete_insert
        +unique_key: trip_id
        +order_by: "(toStartOfMonth(pickup_at), pickup_at, trip_id)"

在模型的 config() 块中

{{ config(
    materialized         = 'incremental',
    engine               = 'ReplacingMergeTree(updated_at)',
    incremental_strategy = 'delete_insert',
    unique_key           = 'trip_id',
    order_by             = '(toStartOfMonth(pickup_at), pickup_at, trip_id)'
) }}

两种方式是等价的。项目级通用模式更适合放在 dbt_project.yml 中;模型级覆盖,或者你希望配置与 SQL 放在一起时,更适合用 config() 块。

关键配置参数

参数它控制什么对应的 ClickHouse 映射
+engine表存储引擎CREATE TABLE 中的 ENGINE = ...
+order_byprimary key / 排序顺序CREATE TABLE 中的 ORDER BY ...;省略时默认为 tuple()
+unique_keydelete_insert 去重所用的键决定插入前删除哪些行
+incremental_strategy增量运行如何更新数据在 ClickHouse 上设为 delete_insert

作用域规则: dbt_project.yml 中的设置从父级向子级层层继承。模型级的 config() 块始终优先于项目配置。把最常用的 engine 设为项目默认值,然后为不同的模型单独覆盖。


3. delete_insert 的工作机制

delete_insert 是 dbt-clickhouse 社区的标准增量策略。它是最接近 Snowflake MERGE INTO 的等价物,但机制不同。

版本要求: delete_insert 使用 ClickHouse 的轻量级删除(lightweight delete),该功能在 22.8 中引入(实验性),并在 23.3+ 中达到生产可用。ClickHouse Cloud 满足这一要求。要启用它,请在 ~/.dbt/profiles.yml 的 ClickHouse target 中加入 use_lw_deletes: true,或在 query_settings 中设置 allow_experimental_lightweight_delete=1。

它做什么

在每次增量运行时:

  1. DELETE 目标表中 unique_key 与本批次任意行匹配的行
  2. INSERT 本批次的所有行
-- Step 1: dbt generates this DELETE
ALTER TABLE analytics.fact_trips
DELETE WHERE trip_id IN (SELECT trip_id FROM incoming_batch);

-- Step 2: dbt generates this INSERT
INSERT INTO analytics.fact_trips
SELECT * FROM incoming_batch;

它与 Snowflake MERGE INTO 的区别

Snowflake 的 merge 策略会生成逐行的 WHEN MATCHED THEN UPDATE / WHEN NOT MATCHED THEN INSERT。ClickHouse 没有 MERGE INTO 语句。delete_insert 通过一次批量删除再加一次完整插入,达到相同的最终结果,每个唯一键一行。

它与 ReplacingMergeTree 的相互作用

delete_insert 是主要的正确性路径。ReplacingMergeTree 是安全网。

如果一次 delete_insert 运行正常完成:表是干净的(每个 trip_id 一行),没有重复。

如果一次 delete_insert 运行中途被打断(DELETE 之后、INSERT 之前崩溃):数据很可能处于无效状态,已被删除的行可能没有被重新插入。下一次成功运行会恢复正确状态,但在失败的 DELETE 与其重跑之间,不要查询该表。

如果某次运行因任何原因产生了重复:ReplacingMergeTree 的后台合并最终会将它们去重,保留版本列取值最大的那一行。

绝不要在没有 delete_insert 的情况下只依赖 RMT,后台合并是异步的,在大表上可能需要数分钟到数小时。

何时使用 append

append 插入新行而不触碰已有行。对于纯插入、行永不更新的表,它是正确的策略,例如不可变的事件日志,或 ID 保证唯一且不会有修正的原始摄取表。append 没有版本要求,也没有 mutation 风险。

对 fact_trips 来说,append 是错误的:一次行程可能事后被修正(车费调整、状态变更),因此同一个 trip_id 会带着新值再次到达。使用 append 时,两个版本会永久累积,而聚合结果(车费的 SUM、行程的 COUNT)会一直高估,直到下一次 RMT 后台合并。只要行可能被更新,就使用 delete_insert。

为什么不用 merge 策略?

merge 策略(delete_insert 之前的历史默认值)会创建一张临时表,用未变更的既有行加上新批次填充它,然后原子地替换原表。与 delete_insert 不同,它不使用轻量级删除,它在每次增量运行时都会重写整张表。对于一张 5000 万行的 fact_trips 表,这代价极高。delete_insert 只处理当前批次中的行;merge 会触及表中的每一行。请使用 delete_insert。


4. FINAL 的放置策略

ReplacingMergeTree 的去重发生在后台,ClickHouse 异步合并 part。在两次合并之间,重复行会共存。FINAL 在读取时强制进行同步去重。

在 dbt 管道中 FINAL 应该放在哪里

放在从 ReplacingMergeTree 源读取并产出干净分析数据的那一层。

对于 NYC Taxi 负载:

trips_raw (RMT)
    ↓
stg_trips (view): SELECT ... FROM trips_raw FINAL   ← FINAL goes here
    ↓
int_trips_enriched (ephemeral CTE)
    ↓
fact_trips (incremental, RMT)                        ← NO FINAL in model
    ↓
Dashboard queries: SELECT ... FROM fact_trips FINAL  ← FINAL goes here (externally)

stg_trips 是 trips_raw 去重的唯一执行点。每个读取 stg_trips 的下游模型都会自动拿到干净、已去重的源数据。你不需要在 int_trips_enriched 或 fact_trips 中使用 FINAL,因为它们读取的是 stg_trips(一个视图,而不是 RMT 表)。

直接从 fact_trips 读取的看板查询和 dbt 测试在外部使用 FINAL。模型本身不内嵌 FINAL,因为那会作用于模型查询内部的每一次扫描,包括从 {{ this }} 读取 max(updated_at) 的 is_incremental() 子查询。

FINAL 对性能的影响

FINAL 带来的延迟与重复行的数量成正比。在一张维护良好的 RMT 表上(后台合并频繁),FINAL 的额外开销很小,因为需要解决的重复很少。而在一张刚刚加载完、有很多未合并 part 的表上,FINAL 会明显更慢。

对于 dbt 测试和校验查询,在 RMT 表上始终使用 FINAL。对于以与 Snowflake 做延迟对比为目的的基准查询,ClickHouse 侧的查询本来就已经使用了 FINAL,因此对比是公平的。


5. generate_schema_name 宏

默认情况下,dbt 会用 profile 中的目标 schema 名给模型 schema 加前缀。如果你的 dbt profile 的目标 schema 是 nyc_taxi_ch,那么带有 +schema: analytics 的模型会落在 nyc_taxi_ch_analytics,而不是 analytics。

这在 Snowflake 中无害(schema 是数据库内的命名空间),但在 ClickHouse 中会产生别扭的名字,因为在那里 schema 就是 数据库。nyc_taxi_ch_analytics 是一个合法的 ClickHouse 数据库名,但它比 analytics 难看,而且与第 3 部分 ClickHouse 架构中使用的目标数据库名不一致。

解决办法是覆盖 generate_schema_name 宏:

-- macros/generate_schema_name.sql
{% macro generate_schema_name(custom_schema_name, node) -%}
  {%- if custom_schema_name is none -%}
    {{ target.schema | lower }}
  {%- else -%}
    {{ custom_schema_name | lower }}
  {%- endif -%}
{%- endmacro %}

这个宏:

  • 当模型指定了 +schema: analytics 时,原样返回 custom_schema_name(转为小写)
  • 对没有自定义 schema 的模型,返回 profile 的目标 schema(转为小写)

| lower 过滤器还确保 schema 名一致地保持小写,以匹配 ClickHouse 区分大小写的标识符规则(Snowflake 第 1 部分使用的是 | upper)。

它放在哪里: macros/generate_schema_name.sql,位于顶层 macros/ 目录中;dbt_project.yml 设置了 macro-paths: ["macros"]。


汇总:NYC Taxi dbt 配置总览

# dbt_project.yml (abbreviated)
models:
  nyc_taxi_dbt_ch:
    staging:
      +schema: staging
      +materialized: view           # no engine — views need none

    intermediate:
      +schema: staging
      +materialized: ephemeral      # inlined as CTE

    analytics:
      +schema: analytics
      +materialized: table
      +engine: "MergeTree()"        # default for dim_* tables

      fact_trips:
        +materialized: incremental
        +engine: "ReplacingMergeTree(updated_at)"
        +incremental_strategy: delete_insert
        +unique_key: trip_id

      agg_hourly_zone_trips:
        +materialized: incremental
        +engine: "ReplacingMergeTree(updated_at)"
        +incremental_strategy: delete_insert
        +unique_key: [hour_bucket, zone_id]
-- stg_trips.sql (staging view — the FINAL enforcement point)
SELECT ... FROM {{ source('raw', 'trips_raw') }} FINAL

-- fact_trips.sql (incremental — no FINAL in model body)
SELECT ... FROM {{ ref('int_trips_enriched') }}
{% if is_incremental() %}
WHERE updated_at > (SELECT max(updated_at) FROM {{ this }})
{% endif %}

本页内容

ZH