- 理解项目中的 view 和表在 ClickHouse 中以何种形式呈现。
- 使用 seed 加载数据,并控制 ClickHouse 类型和表的 layout。
- 为 table 模型配置 ClickHouse engine、sorting key 和 partitioning。
- 将 table 转换为 incremental 模型,并选择合适的 incremental strategy。
- 创建 snapshot。
- 使用 ClickHouse materialized view。
开始之前
请先按照 ClickHouse/jaffle-shop-clickhouse 的 README 操作。其中说明了如何使用 dbt Core 1.x、dbt OSS、dbt v2 或 dbt 平台搭建该项目,如何将其指向本地 ClickHouse (docker) 或 ClickHouse Cloud,如何使用dbt seed 加载示例数据,以及如何执行首次 dbt build。dbt build 成功完成后,再回到本页查看 ClickHouse 特有的示例和配置。
完成 README 中的步骤后,ClickHouse 中应当有两个数据库:
raw:由dbt seed从 CSV file 加载的六张源表 (raw_customers、raw_orders、raw_items、raw_products、raw_stores、raw_supplies) 。jaffle_shop(即你的 profile 中的schema) :六个暂存视图 (stg_*) 和七张 mart 表 (customers、orders、order_items、products、locations、supplies、metricflow_time_spine) 。
schema,请将下面查询中的 jaffle_shop 替换为你的实际取值。
dbt Core 1.x、dbt OSS、dbt v2 与 dbt 平台。 本指南中的所有命令和模型在它们上都完全一致。这些示例已在 dbt Core 1.12 (搭配
dbt-clickhouse 1.10) 和 dbt OSS 2.0 上针对 ClickHouse 26.8 进行过测试;dbt v2 运行的是同一个 adapter,而 dbt 平台运行的是 dbt v2。文中展示的控制台输出来自 dbt Core 1.x,对于少数引擎行为存在差异的地方,文中会特别说明。有关 v2 adapter 的最新状态,请参阅 dbt OSS、dbt v2 与 dbt 平台页面;若要在 dbt 平台上快速上手,请参阅 dbt 文档中的 Connect ClickHouse。clickhouse client、ClickHouse Cloud SQL 控制台或你习惯使用的 SQL 客户端。
项目的物化方式
Jaffle Shop 在dbt_project.yml 中配置其物化类型:staging 模型为视图,marts 为表。
CREATE OR REPLACE VIEW 语句重新构建。它不存储任何数据,因此构建没有任何开销,但每次查询它时,都会针对源表执行该模型的 SQL。ClickHouse 会将模型编译后的 SQL 保存在视图定义中:
INSERT INTO ... SELECT,然后以原子方式将其与上一版本互换。其查询性能远优于 view,代价是占用存储空间,且每次都要重建整张表。来看看 dbt 为 orders 这个 mart 创建的表:
MergeTree;它也没有声明sorting key,因此adapter使用 ORDER BY tuple(),也就是数据完全不排序。对于示例项目来说这没什么问题,但在真实的表中,你应当明确指定这两项,这也正是后续几节要做的事情。物化类型页面列出了adapter支持的所有表配置。
使用 seeds 加载数据
Jaffle Shop 使用 dbt seeds 从seeds/jaffle-data 目录下的 CSV 文件加载原始数据。Seeds 的用途是小型静态参考数据 (代码表、映射关系) ,而非向数据仓库批量加载数据;本项目只是为了方便才使用它,让你无需额外的摄取工具即可上手——这也是为什么除非传入 --vars '{"load_source_data": true}',否则 seeds 默认处于禁用状态。
即便如此,seeds 仍是了解 dbt 如何创建 ClickHouse 表的好切入点。dbt 会为每个 CSV 列推断列类型,而不同引擎推断出的类型有所差异:
当类型很重要时,请用
column_types 显式固定类型。本项目已在 dbt_project.yml 中对 raw_stores seed 的 opened_at 列这样处理:
engine、order_by 和 partition_by。例如,若要让 raw_orders seed 按下单时间排序并按月分区,可在 CSV 文件旁添加一个属性文件 seeds/jaffle-data/_raw_orders.yml:
这些 ClickHouse seed 配置应写在 properties 文件中,而不要放在
dbt_project.yml 的 seeds: 下,通过 +order_by 或 +engine 键来指定。dbt Core 1.x 两种写法都接受,但 dbt v2 只识别 properties 文件中的配置,并会拒绝 dbt_project.yml 中的这些键,报错 Unrecognized key ... Custom keys must go under +meta。dbt seed --full-refresh 会删除并重新创建该表,因此请先运行它,然后再构建任何直接依赖该表数据的对象 (例如本指南后文中的 materialized view) 。
为 ClickHouse 配置表
orders mart 是最自然的起点:它会被 customers mart 以及项目的 metrics 查询,而且它是一张带 timestamp 的事件型表。在 models/marts/orders.sql 顶部添加一个 config 块,用于指定 engine、sorting key 和分区方案:
materialized='table' 重复了 dbt_project.yml 中已为 marts 设定的内容,这样在之后将该模型切换为增量模式时,模型本身仍然是自描述的。仅重新构建该模型:
engine、order_by 和 partition_by 之外,table 模型还支持 primary_key、ttl、settings、query_settings、projections 和 indexes,而列可以通过模型契约 (model contract) 指定 codec 和 ttl。这些配置均在物化类型页面中有详细说明。
创建增量模型
对于 62,000 行数据来说,每次运行都从头重建orders 并无不可,但若表每天新增数百万行,就行不通了。dbt 的增量物化只处理自上次运行以来发生变化的行。将 orders 模型改为增量模型需要两处补充:
unique_key:用于标识一行的列,这里是order_id。adapter 会用它替换被再次处理的行,而不是产生重复数据。- 增量过滤器:包裹在
{% if is_incremental() %}中的where子句,只选取需要处理的行。它在增量运行时生效,而在表首次构建 (或使用--full-refresh重建) 时不生效。订单带有 timestamp,因此该过滤器会将ordered_at与表中已有的最新值进行比较,该值通过{{ this }}变量引用。
models/marts/orders.sql,使 config 块和模型末尾部分如下所示:
stg_orders 会将 ordered_at 截断到天,因此过滤器使用 >=:每次运行都会重新处理最新一整天的数据,并且借助 unique_key,已加载的行会被替换而不是重复写入。正因如此,对于同一天稍晚到达的订单也能安全处理。
运行该模型。由于表已经存在,这第一次运行其实就是一次增量运行:只会重新处理最新的一天。
nutellaphone who dis? jaffle,税率为 Philadelphia 的 6%,因此项目的数据测试仍然通过。运行整个项目,让暂存视图和 order_items 表先于 orders 看到这些新行:
customers 数据集市也已包含这位新客户:
内部实现
ClickHouse 的 query log 中可以看到 adapter 为执行增量更新所运行的语句:- 创建表
orders__dbt_new_data,并将模型 SQL (含增量过滤器) 的查询结果 insert 到该表中。在上面的运行中共写入 378 行:已加载的最近一天的 377 个订单,加上新增的那一个。 - 创建一个与
orders结构相同的表orders__dbt_tmp,并把orders中order_id未出现在orders__dbt_new_data里的所有行复制进去。 - 将
orders__dbt_new_data的所有行 insert 到orders__dbt_tmp。正是步骤 2 和 3 实现了替换最近一天的行,而非重复写入。 - drop 掉
orders__dbt_new_data。 - 通过原子的
EXCHANGE TABLES语句将orders__dbt_tmp与orders进行 swap (中间先 rename 为orders__dbt_backup) ,此时orders中保存的即为新版本。 - drop 掉旧版本。
追加策略
append 策略会将模型选出的行直接插入到 target table 中。它不会创建 temporary table,也不会复制任何数据,因此是增量运行中开销最低的方式。代价是它同样不做任何去重:如果增量过滤器选中了表中已有的行,该行就会重复出现两次。因此请仅将它用于不可变的事件类数据,并确保 filter 只会选中真正的新行。
在 ordered_at 已按天截断的情况下,这意味着要把 filter 改为 >。修改模型:
orders 的语句只有一条 INSERT INTO jaffle_shop.orders ... SELECT ...,其中包含该模型的 SQL 以及增量过滤器,并且只写入了一行。
Delete and insert 策略
一直以来,ClickHouse 对更新和删除的支持都比较有限,只能通过异步 变更 来实现。这类操作往往会带来极高的 IO 开销,通常应尽量避免。ClickHouse 22.8 引入了 轻量级删除,ClickHouse 25.7 引入了 轻量级更新。有了这些特性,虽然数据是以异步方式 materialized 的,但从用户的角度看,单条删除或更新语句的效果会立即可见。delete+insert 策略基于轻量级删除实现,通过 incremental_strategy 参数进行配置:
- 创建一个 temporary table (
orders__dbt_new_data_<run_id>) ,并将 model 选出的行 insert 到其中。 - 针对 temporary table 中出现的每个
order_id,对orders执行DELETE。 - 将 temporary table 中的行 insert 到
orders。 - drop 该 temporary table。
Insert overwrite 策略 (experimental)
insert_overwrite 策略会整分区替换,因此需要配置 partition_by,例如 orders 上按月分区的配置。其执行步骤如下:
- 创建一个与
orders结构相同的暂存表 (orders__dbt_new_data_<run_id>) 。 - 仅将模型选出的行插入暂存表。
- 从
system.parts中列出暂存表内存在的分区。 - 使用
ALTER TABLE ... REPLACE PARTITION ... FROM,用暂存表精确替换orders中的这些分区。 - 删除暂存表。
- 比默认策略更快,因为无需复制整张表。
- 比其他策略更安全,因为在 INSERT 操作成功完成之前不会修改原始表:若中途失败,原始表保持不变。
- 实现了“分区不可变性”这一数据工程最佳实践,从而简化增量与并行数据处理、回滚等操作。
microbatch 策略和 on_schema_change。
创建 snapshot
dbt snapshots 用于记录可变表中的行随时间发生的变化,使分析人员能够回溯查看过去任意时间点的数据状态。它们实现了类型 2 缓慢变化维:每个版本的行都会连同其有效的时间间隔一并存储。customers 数据集市就很合适:每当客户再次下单,count_lifetime_orders、lifetime_spend 和 customer_type 都会随之变化。在继续之前,请将 orders 模型改回增量章节中的默认增量策略 (移除 incremental_strategy='append',并将过滤条件改回 >=) ,这样当天稍后产生的订单才能被采集到。
自 dbt 1.9 起,snapshots 通过 YAML 定义。创建 snapshots/customers_snapshot.yml:
check 策略会在每次运行时比较 current snapshot 与 source 中所列出的列,只要其中任意一列发生变化,就记录一个新版本。如果你的模型中有一个可靠的 “last updated” 时间戳列,则 timestamp 策略开销更低:设置 strategy: timestamp 和 updated_at: <column>。Jaffle Shop 的 last_ordered_at 被截断到天,因此无法捕获同一天内的第二笔订单,这正是本示例使用 check 的原因。
创建第一个 snapshot:
generate_schema_name macro 会把所有 relation 都放到非 production target 的 target schema 中,因此 snapshot 上的 schema 配置只在使用 prod target 时才生效。该表为每个客户保存一行,并包含 dbt 用于记账的列 dbt_valid_from 和 dbt_valid_to;对于某一行的当前版本,后者为 NULL:
orders 和 customers 反映这笔新订单,然后创建第二个快照:
dbt_valid_to 被关闭,而新版本 (此时已是拥有两笔订单的 returning 客户) 处于打开状态。Danny 没有变化,因此他的行保持原样:
customers_snapshot__snapshot_upsert 中构建新版本的 snapshot,然后通过 EXCHANGE TABLES 将其替换上线 (若 server 不支持 exchange tables,则改用 drop 加 rename 的方式) ,因此读取器看到的要么是旧版本的 snapshot,要么是新版本,不会出现中间状态。配置参考请参见物化类型页面的 snapshot 部分。
使用 materialized view
到目前为止,所有内容都需要执行dbt run 才能将新数据加载到模型中。ClickHouse materialized view 的工作方式则不同:它们本质上是插入触发器。每当有行块插入源表,视图的 SELECT 就会对其进行转换并写入目标表,无需任何调度。adapter 通过 materialized_view 物化暴露这一能力。
创建 models/marts/daily_store_revenue.sql,直接从原始订单表读取数据,统计每个门店每天的订单数和收入:
engine 和 order_by 作用于目标表。SummingMergeTree 在合并 parts 时,会将 sorting key 相同的行的数值列相加,这正是按天、按门店聚合所需要的行为。
_mv 后缀的 materialized view 本身,后者通过 TO clause 指向该 target table。默认情况下 (catchup=True) ,target table 还会用已有的订单数据完成 backfill:
sum() 和 GROUP BY 进行聚合:SummingMergeTree 只在后台合并 parts 时才会折叠相同键的行,在此之前,两笔 Brooklyn 订单在表中仍是两行。因此,使用求和类和聚合类引擎时,务必在读取时聚合 (或使用 FINAL) 。与此同时,在下一次 dbt run 之前,orders 增量模型中 Danny 仍然只有一笔订单。
后续执行 dbt run 会保留target table及其数据,只更新 view definition;如果变更允许,会通过 ALTER TABLE ... MODIFY QUERY 完成,因此把该模型保留在项目中是安全的。dbt run --full-refresh 则会重建target table并重新 backfill (除非 catchup 为 False) 。其余内容可参见 materialized views 页面:使用 on_schema_change 处理 schema 变更、通过 catchup 禁用 backfill、可刷新 materialized views、多个 view 写入同一target table,以及将target table定义为独立模型。