04 重建 dbt 流水线
用 dbt-clickhouse 在 ClickHouse 上重建 Medallion 流水线,delete_insert 增量模型、ReplacingMergeTree、可刷新的 materialized view,并创建区域字典。
起点
模块 03 已完成:ClickHouse Cloud 服务已上线且可访问,
setup.sh 已把带有 CLICKHOUSE_HOST 和
CLICKHOUSE_PORT 的 .clickhouse_state 写到磁盘上。所有目标表和 staging 视图都已存在,default.trips_raw、
两个 staging 视图(stg_trips、stg_taxi_zones)、六张 analytics 表
(fact_trips、agg_hourly_zone_trips、dim_taxi_zones、dim_payment_type、dim_vendor、
dim_date),以及可刷新的 materialized view analytics.mv_live_trip_feed,模块
03 的 dbt run 已经把这七个全部建好。除了
mv_live_trip_feed(它已装着那次构建产生的一行快照数据),每张 analytics 表都还是空的。除此之外只有
default.trips_raw 有数据:大约 5000 万行。Snowflake 生产者
仍在运行,因此 ClickHouse 落后 Snowflake 大约一个迁移窗口的
时长。预留大约 30 分钟。
为什么
模块 03 证明了 ClickHouse 能装下 5000 万行。它并没有证明流水线可以 跑在 ClickHouse 上,staging 视图、增量事实表、维度表重载, 以及那些在合作伙伴看到之前就抓住坏模型的测试。这就是本 模块要重建的东西:与模块 01 相同的那批 Medallion 模型,改用 dbt-clickhouse 而不是 dbt-snowflake 来表达,针对模块 03 已经创建好的那些表运行。
模型逻辑本身没有任何变化,stg_trips 依然做类型转换并提取 JSON、
int_trips_enriched 依然关联各维度、fact_trips 依然最终做到每个行程
一行。变化的是底下的物化层:没有 MERGE INTO、没有 Snowflake
Task、没有 cluster_by。这个模块要说明的是,这次迁移不是一次性的数据
搬运,合作伙伴团队每天按同样的排期运行、由同一次 dbt test 把关的
那条流水线,在它底下的数仓换掉之后依然照常工作。
概念:底层原理
这套模型与 Snowflake 流水线有四处不同。完整的配置
参考见
ClickHouse 上的 dbt,下面是你在
步骤 1 运行 dbt run 之前所需的精简版。关于这些模型所替代的源端流水线,
请见
Snowflake 上的 dbt。
1. delete_insert 取代 MERGE。 ClickHouse 没有 MERGE INTO 语句。Snowflake
流水线用 incremental_strategy: merge 来 upsert fact_trips 和
agg_hourly_zone_trips 的地方,ClickHouse 模型改用 incremental_strategy: delete_insert:
dbt 先删除与传入批次的 unique_key 匹配的行,然后插入该批次。
对 fact_trips 来说,unique_key 是 trip_id,而增量过滤条件的水位线取
updated_at 而不是 pickup_at,一次车费修正会以相同的 pickup_at、更新的
updated_at 重新插入同一个 trip_id,所以用 pickup_at 做水位线会
无声地漏掉它。
2. ReplacingMergeTree 是 delete_insert 底下的安全网,而不是它的
替代品。 两个增量模型都声明为 ReplacingMergeTree(updated_at)。如果一次
delete_insert 正常完成,表里已经是每个键一行,引擎没有什么需要清理的。如果一次运行
中途被打断,在删除之后、插入之前崩溃,后台合并最终会对任何残留行去重,
保留 updated_at 最大的那一行。绝不要指望仅靠 ReplacingMergeTree 来完成
delete_insert 本应做的去重工作:后台合并是异步的,
在这种规模的表上可能滞后数分钟到数小时。
3. 可刷新的 materialized view 取代定时任务。 Snowflake 流水线
用一个运行存储过程的定时 Task 来保持滚动聚合最新。
ClickHouse 的 dbt 项目则把 mv_live_trip_feed 声明为
materialized = 'materialized_view' 加 engine = 'ReplacingMergeTree(refreshed_at)',
由 dbt run 构建为一个可刷新的 materialized view,而在 Snowflake 上达到同样效果需要
30 多行 CREATE TASK DDL。本模块不会开启刷新间隔;
要开启的话,是一条本实验没有脚本化的手动
ALTER TABLE analytics.mv_live_trip_feed MODIFY REFRESH EVERY ... 语句(原因见模块 05)。
4. mv_live_trip_feed 在 Snowflake 侧没有对应实现。 它不是对现有模型的转换,
而是本次迁移新增的能力。标准的 ClickHouse
materialized view 在每次 INSERT 时触发一次,而且只能看到那一批次里的数据行,因此
它无法正确计算像总行程数或平均车费这样的全生命周期聚合。而
REFRESHABLE materialized view 会按排期重跑它的整个查询,这里是
SELECT ... FROM {{ ref('fact_trips') }},所以每次刷新都能看到完整的
表。本课程的 Snowflake 环境从来没有这个选项。
dbt 实际构建的模型,以及每一个何时获得数据:
| 模型 | 层 | 物化方式 | 何时填充数据 | 说明 |
|---|---|---|---|---|
stg_trips | staging | 视图 | 步骤 1(每次运行) | 类型转换,对 trip_metadata 使用 JSONExtract* |
stg_taxi_zones | staging | 视图 | 步骤 1(每次运行) | 区域维度直通 |
int_trips_enriched | staging | Ephemeral | ,(内联为 CTE) | 所有维度关联;没有物理表 |
fact_trips | analytics | 增量 | 步骤 1 | 以 trip_id 为键的 delete_insert,水位线取 updated_at;ORDER BY (toStartOfMonth(pickup_at), pickup_at, trip_id) |
agg_hourly_zone_trips | analytics | 增量 | 模块 05,切换之后 | 滚动 2 小时窗口;增量过滤条件只匹配实时生产者写入的行,见下方步骤 1 |
dim_taxi_zones | analytics | 表 | 步骤 1 | 每次运行整体重载;步骤 2 中区域字典的数据来源 |
dim_payment_type | analytics | 表 | 步骤 1 | 每次运行整体重载 |
dim_vendor | analytics | 表 | 步骤 1 | 每次运行整体重载 |
dim_date | analytics | 表 | 步骤 1 | 静态日期轴,2009-2029;每次运行整体重载 |
mv_live_trip_feed | analytics | Materialized view(可刷新) | 模块 03 的 dbt run;刷新间隔从未开启 | 没有 Snowflake 对应物,见上文第 4 点 |
步骤 1:填充分析层
激活你在模块 00 构建的 dbt-clickhouse venv,然后第二次运行
dbt run。模块 03 已经对着空表跑过一次以创建 schema;这次运行
背后有真实数据,trips_raw 现在装着 5000 万行。
cd "$(git rev-parse --show-toplevel)/workshop_public/snowflake_migration_lab/03-migrate-to-clickhouse"
source .venv/bin/activate
source .env && source .clickhouse_state
cd "$(git rev-parse --show-toplevel)/workshop_public/snowflake_migration_lab/03-migrate-to-clickhouse/dbt/nyc_taxi_dbt_ch"
dbt rundbt 从 ~/.dbt/profiles.yml 读取 ClickHouse 连接信息,其模板为
dbt/nyc_taxi_dbt_ch/profiles.yml.example,也就是模块 03 用来创建
空 schema 的那个 profile。
预期: 大约 8-12 分钟(增量模型要处理 5000 万 行)。
这次运行之后 agg_hourly_zone_trips 会是空的,这是预期结果,不是
故障。 它的增量过滤条件是 WHERE pickup_at >= now() - INTERVAL 2 HOUR,
只会匹配实时生产者写入的行。你刚迁移过来的每一行都是
历史数据,所以没有一行落在以当下为起点的 2 小时窗口内。这张表
会一直保持为空,直到模块 05 的切换启动 ClickHouse 生产者,不要花时间
把它当成坏掉的流水线来排查。
然后运行测试套件:
dbt test预期: 所有测试通过。
验证:
SELECT formatReadableQuantity(count()) FROM analytics.fact_trips FINAL;
-- Expected: ~50 million
SELECT count() FROM analytics.agg_hourly_zone_trips;
-- Expected: 0 (normal — populated after cutover in module 05)
SELECT count() FROM analytics.dim_taxi_zones;
-- Expected: 265步骤 2:创建区域字典
analytics.dim_taxi_zones 现在已装入全部 265 个 NYC TLC 区域。构建
analytics.taxi_zones_dict,一个以那张表为后端的内存字典,这样下游
查询就能用 dictGet() 查出某个区域所属的行政区,而不必用 JOIN。
字典相比关联能带来什么。 字典只加载进内存一次,然后一直保持
热态;之后每次查询对它做查找实际上是零成本的。而对
dim_taxi_zones 做 JOIN,每次运行都要重新读取并重新匹配这张维度表。对于像这样
一张小而极少变化的参考表,265 行,每次
dbt run 都整体重载,这个取舍毫无悬念。模块 05 的基准测试会直接用 dictGet 查询
taxi_zones_dict,因此这一步是那个模块的硬依赖,
不是可选的附加项。
先 source 连接信息,然后通过
clickhouse-client 或 HTTP API 加载字典 DDL,用哪种取决于你手头有什么:
cd "$(git rev-parse --show-toplevel)/workshop_public/snowflake_migration_lab/03-migrate-to-clickhouse"
source .clickhouse_state
# Via clickhouse-client
clickhouse-client \
--host "${CLICKHOUSE_HOST}" \
--port 9440 \
--user default \
--password "${CLICKHOUSE_PASSWORD}" \
--secure \
--multiquery \
< scripts/04_create_dictionary.sql
# Or via HTTP API
curl "https://${CLICKHOUSE_HOST}:${CLICKHOUSE_PORT}/" \
--user "default:${CLICKHOUSE_PASSWORD}" \
--data-binary @scripts/04_create_dictionary.sql验证:
-- Should return 'Manhattan' for zone 42
SELECT dictGet('analytics.taxi_zones_dict', 'borough', toUInt16(42));
-- Should show status = LOADED, element_count = 265
SELECT name, status, element_count
FROM system.dictionaries
WHERE name = 'taxi_zones_dict';如何确认已完成
SELECT formatReadableQuantity(count()) FROM analytics.fact_trips FINAL;
-- Expected: ~50 millionSELECT count() FROM analytics.agg_hourly_zone_trips;
-- Expected: 0这里的 0 是正确的,不是故障。agg_hourly_zone_trips 的增量过滤条件
(WHERE pickup_at >= now() - INTERVAL 2 HOUR)只匹配实时生产者写入的
行,而此刻 ClickHouse 中的每一行都是模块 03 的迁移脚本搬来的历史
数据,相对 now(),没有一行的年龄小于 2 小时。这张表
只会在模块 05 的切换启动 ClickHouse 生产者之后才填充数据;在那之前,任何
由这张表支撑的仪表板图表都会显示无数据,而这是预期结果,
不是需要排查的问题。
SELECT count() FROM analytics.dim_taxi_zones;
-- Expected: 265cd "$(git rev-parse --show-toplevel)/workshop_public/snowflake_migration_lab/03-migrate-to-clickhouse/dbt/nyc_taxi_dbt_ch"
dbt test预期:所有测试通过。
-- Should return a borough name, e.g. 'Manhattan'
SELECT dictGet('analytics.taxi_zones_dict', 'borough', toUInt16(42));结束状态
分析层已填充数据并通过测试:fact_trips 装着大约 5000 万行,
dim_taxi_zones、dim_payment_type、dim_vendor 和 dim_date 全部加载完毕,
dbt test 端到端通过,analytics.taxi_zones_dict 已上线并能通过 dictGet()
返回行政区。agg_hourly_zone_trips 仍然是空的,这是设计如此,不是
缺陷,并且会一直如此,直到模块 05 的切换。
仪表板、ClickHouse 与 Snowflake 的基准对比以及切换,属于模块 05,不是这 一个。
Snowflake 生产者仍在运行,Snowflake 与 ClickHouse 之间的间隔也依然 存在。 本模块中没有任何操作触碰过生产者或迁移脚本, 模块 05 会刻意弥合那个间隔,用的正是模块 03 预告过的那个受控两遍 步骤。现在不要停掉生产者。