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 支持五种物化方式。选择依据是更新模式,而不是个人偏好。
| 物化方式 | 物理对象 | 何时使用 |
|---|---|---|
view | ClickHouse 视图 | Staging 模型:清洗并转换源数据类型;无存储成本;每次查询时重新计算 |
ephemeral | 无对象(内联为 CTE) | 通过 JOIN 组合多个 staging 模型的中间模型;避免创建冗余的物理表 |
table | 先在一个临时关系中构建完整的替换版本,再通过 EXCHANGE TABLES(或在旧版本上用成对重命名)原子地换入;换入后删除旧表 | 每次 dbt run 都完全替换的小型维度表;不需要局部更新。注意: 对大表来说完全重建不可行,任何超过几千行的表都应使用 incremental。 |
incremental | 首次运行时 CREATE TABLE;后续运行采用选择性 UPDATE 模式 | 事实表和预聚合表,每次运行只需处理新增/变更的行 |
materialized_view | ClickHouse 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_by | primary key / 排序顺序 | CREATE TABLE 中的 ORDER BY ...;省略时默认为 tuple() |
+unique_key | delete_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。
它做什么
在每次增量运行时:
- DELETE 目标表中
unique_key与本批次任意行匹配的行 - 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 %}