- 创建 dbt 项目并设置 ClickHouse 适配器。
- 定义模型。
- 更新模型。
- 创建增量模型。
- 创建快照模型。
- 使用 materialized view。
设置
准备 ClickHouse
表
roles 中的 created_at 列默认值为 now()。稍后我们会用它来识别模型的增量更新——请参见增量模型。s3 函数从公共端点读取源数据,并将数据插入表中。运行以下命令来填充这些表:
连接到 ClickHouse
-
创建一个 dbt 项目。在本例中,我们以
imdbsource 为项目命名。出现提示时,选择clickhouse作为数据库 source。 -
使用
cd进入项目目录: - 此时,你需要使用自己选择的文本编辑器。在下面的示例中,我们使用常见的 VS Code。打开 IMDB 目录后,你应该会看到一组 yml 和 sql 文件:
-
更新你的
dbt_project.yml文件,指定第一个模型actor_summary,并将 profile 设为clickhouse_imdb。 -
接下来,我们需要向 dbt 提供 ClickHouse 实例的 connection details。将以下内容添加到
~/.dbt/profiles.yml中。请注意,你需要修改 user 和 password。有关其他可用设置的说明,请参见这里。 -
在 IMDB 目录中,执行
dbt debug命令,确认 dbt 是否能够连接到 ClickHouse。确认响应中包含Connection test: [OK connection ok],表示连接成功。
创建简单的视图物化
CREATE VIEW AS 语句重建为视图。这样无需额外存储数据,但查询速度会比表物化类型慢。
-
在
imdb文件夹下,删除目录models/example: -
在
models文件夹中的actors目录下创建一个新文件。这里创建的每个文件都对应一个 actor 模型: -
在
models/actors文件夹中创建schema.yml和actor_summary.sql这两个文件。文件schema.yml定义了我们的表。之后,这些表就可以在 macro 中使用。编辑models/actors/schema.yml,使其包含以下内容:actors_summary.sql定义了实际的模型。请注意,在config函数中,我们还指定将该模型在 ClickHouse 中 materialize 为视图。我们的表是通过schema.yml文件中的source函数引用的,例如source('imdb', 'movies')指向imdbdatabase 中的movies表。将models/actors/actors_summary.sql编辑为以下内容:请注意,我们在最终的 actor_summary 中加入了updated_at列。后续会将其用于增量物化。 -
在
imdb目录下执行命令dbt run。 -
dbt 会按要求将该模型在 ClickHouse 中表示为一个视图。现在,我们可以直接查询该视图。该视图会创建在
imdb_dbtdatabase 中——这是由clickhouse_imdbprofile 下~/.dbt/profiles.yml文件中的 schema parameter 决定的。通过查询这个视图,我们可以用更简单的语法复现先前查询的结果:
创建表物化
INSERT TO SELECT。请注意,这张表每次都会被重新构建,也就是说,它不是增量式的。因此,较大的结果集可能会导致较长的执行时间——请参阅 dbt Limitations。
-
修改文件
actors_summary.sql,将materialized参数设置为table。注意ORDER BY的定义方式,以及这里使用的是MergeTree表引擎: -
在
imdb目录中执行命令dbt run。此次执行可能会稍慢一些——在大多数机器上大约需要 10 秒。 -
确认表
imdb_dbt.actor_summary已创建:你应该会看到包含相应数据类型的表: -
确认该表返回的结果与之前的结果一致。注意,现在模型已物化为表,响应时间有了明显改善:
你也可以继续对此模型执行其他查询。例如,出场次数超过 5 次的演员中,哪些演员参演的电影平均评分最高?
创建增量物化
-
首先,我们将模型改为
incremental类型。此更改需要:- unique_key - 为确保适配器能够唯一标识各行,我们必须提供一个 unique_key——在本例中,查询中的
id字段就足够了。这样可以确保物化后的表中不会出现重复行。有关唯一性约束的更多信息,请参见这里。 - Incremental filter - 我们还需要告诉 dbt,在增量运行时应如何识别哪些行发生了变化。这可以通过提供一个增量表达式来实现。对于事件数据,这通常会涉及一个 timestamp;因此这里使用的是
updated_attimestamp 字段。该列在插入行时默认值为 now(),从而可以识别新增的角色。此外,我们还需要识别另一种情况,即新增了 actor。使用{{this}}变量表示现有的物化表后,就得到这个表达式:where id > (select max(id) from {{ this }}) or updated_at > (select max(updated_at) from {{this}})。我们将它嵌入{% if is_incremental() %}条件中,以确保它只在增量运行时使用,而不会在首次构建表时使用。有关为增量模型过滤行的更多信息,请参见 dbt 文档中的这段讨论。
actor_summary.sql:请注意,我们的模型只会处理roles和actors表中的更新和新增数据。若要覆盖所有表,建议将此模型拆分为多个子模型——每个子模型都有各自的增量条件。随后,这些模型可以相互引用并关联起来。有关模型间交叉引用的更多信息,请参见此处。 - unique_key - 为确保适配器能够唯一标识各行,我们必须提供一个 unique_key——在本例中,查询中的
-
执行
dbt run,并确认生成表中的结果: -
现在,我们将向模型添加数据,以演示增量更新。将我们的演员 “Clicky McClickHouse” 添加到
actors表中: -
让“Clicky”出演 910 部随机电影:
-
通过查询底层源表并绕过所有 dbt 模型,确认他如今确实已是出场次数最多的演员:
-
运行一次
dbt run,并确认我们的模型已更新,且与上述结果一致:
内部原理
- 适配器 会创建一个临时表
actor_sumary__dbt_tmp。发生变化的行会被流式写入该表。 - 接着会创建一个新表
actor_summary_new,。随后,旧表中的行会从旧表流式传输到新表,同时检查这些行的 ID 是否不存在于临时表中。这样可以有效处理更新和重复数据。 - 临时表中的结果会被流式传输到新的
actor_summary表中: - 最后,通过
EXCHANGE TABLES语句以原子方式将新表与旧版本交换。随后再删除旧表和临时表。
追加策略 (仅插入模式)
incremental_strategy。可将其设置为 append。设置后,更新的行会直接插入目标表 (即 imdb_dbt.actor_summary) ,不会创建临时表。
注意:仅追加模式要求数据是不可变的,或者可以接受重复数据。如果你需要支持已修改行的增量表模型,请不要使用此模式!
为了演示此模式,我们将再添加一位新演员,并在 incremental_strategy='append' 的情况下重新执行 dbt run。
-
在 actor_summary.sql 中配置仅追加模式:
-
再添加一位著名演员 —— Danny DeBito
-
让 Danny 参演 920 部随机电影。
-
执行一次
dbt run,并确认 Danny 已添加到 actor_summary 表中
imdb_dbt.actor_summary 表中,不会创建表。
删除和插入模式 (Experimental)
incremental_strategy 参数为模型配置此模式,即
- 适配器会创建一个临时表
actor_sumary__dbt_tmp。发生变更的行会被流式写入该表。 - 对当前的
actor_summary表执行一条DELETE。根据actor_sumary__dbt_tmp中的 id 删除对应的行。 - 使用
INSERT INTO actor_summary SELECT * FROM actor_sumary__dbt_tmp将actor_sumary__dbt_tmp中的行插入actor_summary。
insert_overwrite 模式 (Experimental)
- 创建一个与增量模型 relation 结构相同的暂存 (临时) 表:
CREATE TABLE {staging} AS {target}。 - 仅将新记录 (由 SELECT 生成) 插入暂存表。
- 仅将新分区 (即暂存表中存在的分区) 替换到目标表中。
这种方法有以下优点:
- 它比默认策略更快,因为无需复制整个表。
- 它比其他策略更安全,因为在 INSERT 操作成功完成之前,不会修改原始表:如果中途失败,原始表不会被修改。
- 它实现了数据工程中“分区不可变性”的最佳实践,从而简化增量和并行数据处理、回滚等操作。
创建快照
-
在 snapshots 目录中创建一个
actor_summary文件。 -
将 actor_summary.sql 文件的内容更新为以下内容:
select查询定义了你希望随时间推移进行快照的结果。ref函数用于引用我们之前创建的 actor_summary 模型。- 我们需要一个时间戳列来标识记录变更。这里可以使用
updated_at列 (参见创建增量表模型) 。strategy参数表示我们使用时间戳来标记更新,而updated_at参数则指定使用哪一列。如果你的模型中没有这个列,也可以改用 check 策略。这种方式效率会低很多,并且需要用户指定要比较的列列表。dbt 会比较这些列的当前值和历史值,并记录所有变化 (如果值相同,则不执行任何操作) 。
-
运行命令
dbt snapshot。
snapshots DB 中已创建名为 actor_summary_snapshot 的表 (由 target_schema parameter 决定) 。
-
对这些数据进行抽样后,你会看到 dbt 添加了 dbt_valid_from 和 dbt_valid_to 这两列。后者的值为 null。后续运行会更新这一点。
-
让我们最喜欢的演员 Clicky McClickHouse 再出演 10 部电影。
-
在
imdb目录中重新运行 dbt run 命令。这将更新增量模型。完成后,运行 dbt snapshot 以捕获这些变更。 -
如果我们现在查询这个快照,会发现 Clicky McClickHouse 有 2 行。我们之前的记录现在有了 dbt_valid_to 值。新记录在 dbt_valid_from 列中的值与其相同,而 dbt_valid_to 的值为 null。如果存在新行,这些行也会被追加到快照中。
使用 seed
-
我们先从现有数据集中生成一份类型代码列表。在 dbt 目录中,使用
clickhouse-client创建文件seeds/genre_codes.csv: -
执行
dbt seed命令。这会在数据库imdb_dbt中创建一个新表genre_codes(由 schema 配置定义) ,并将 csv 文件中的行加载到该表中。 -
确认这些数据已加载: