Snowflake MigrationClickHouse Workshops

05 基准测试与切换

在 ClickHouse 上重建仪表板、对两个引擎跑完全部七个查询的基准测试、切换生产者、验证一致性,并拆除资源。

起点

模块 04 已完成:分析层已填充数据并通过测试。analytics.fact_trips 装着大约 5000 万行;analytics.dim_taxi_zones、analytics.dim_payment_type、 analytics.dim_vendor 和 analytics.dim_date 全部加载完毕;dbt test 端到端 通过;analytics.taxi_zones_dict 已上线并能通过 dictGet() 返回行政区。analytics.agg_hourly_zone_trips 仍然是空的,这是设计如此,不是 缺陷,并且会一直如此,直到本模块的切换。Snowflake 生产者仍在 运行,Snowflake 与 ClickHouse 之间的间隔也依然存在。预留大约 45 分钟。

为什么

在此之前的每个模块都是准备工作。模块 03 证明了 ClickHouse 能装下 5000 万行;模块 04 证明了 dbt 流水线能对着它运行。这两件事,单独看, 都不足以让合作伙伴签字确认"已迁移",那还需要两样东西: 一个数字,和一次切换。

数字来自步骤 2 的基准测试:你迁移方案里那同样的七个查询, 对 Snowflake 和 ClickHouse 背靠背运行,各跑三次取中位数。 正是这一步,把"ClickHouse 应该更快"变成一个具体、可论证的提速倍数,让合作伙伴 能拿到自己的利益相关方面前去讲。

步骤 3 的切换是另一半。到目前为止每个模块都让两套系统 并行运行,Snowflake 是记录系统,ClickHouse 在它后面追赶。 一次从未真正搬动写入路径的迁移,是一次拷贝,不是一次迁移。步骤 3 会停掉 Snowflake 生产者、弥合它自模块 01 起持续写入所一直保持着的那个间隔, 并转而开始把新行程写入 ClickHouse,那一刻 ClickHouse 成为记录系统。

这也是模块 06 书面考核之前的最后一个学员模块,而那场考核是 针对你在这里产出的东西开卷进行的。仪表板、基准测试 CSV 和 一致性检查都必须是真实完成的,然后你才能在步骤 5 拆除任何东西。

概念:底层原理

BI 层的机制。 步骤 1 中 bash superset/add_clickhouse_connection.sh 直接调用 Superset REST API,不需要在 Superset 界面里手动点选。它会 注册 ClickHouse 连接,然后导入已提交的仪表板导出包, 在模块 01 建好的三个 Snowflake 仪表板之外再加上四个 ClickHouse 仪表板 (共七个):

仪表板对应它演示了什么
CH,Operations Command CenterSnowflake 仪表板 1实时 fact_trips 数据(切换后);相同的 KPI,更快的查询
CH,Executive Weekly ReportSnowflake 仪表板 2把 QUALIFY 改写为 ROW_NUMBER() 子查询
CH,Driver & Quality AnalyticsSnowflake 仪表板 3用 JSONExtractString 取代 Snowflake 的 LATERAL FLATTEN
CH,Capabilities Showcase(新增,没有 Snowflake 对应物)近似函数、字典关联、SAMPLE 子句

这七个基准查询考察什么。 步骤 2 会用你迁移方案里同样的七个查询 对两个引擎运行并比较实际耗时。每一个都针对模块 02 方案中的 某个特定方言差异或引擎特性:

查询考察内容
Q1按行政区的小时级营收
Q2滚动 7 日平均行程距离
Q3Top-10 行程,Snowflake QUALIFY 对比 ClickHouse ROW_NUMBER() 子查询
Q4司机评分,Snowflake LATERAL FLATTEN 对比 ClickHouse JSONExtractString
Q5高峰计价,Snowflake VARIANT 对比 ClickHouse String + JSONExtract*
Q6小时级聚合,Snowflake MERGE 对比 ClickHouse ReplacingMergeTree
Q7CDC / 实时数据新鲜度

切换间隔。 Snowflake 生产者自模块 01 起就在以约每分钟 60 条行程的速度 写入,而且从未停过。模块 03 的迁移脚本捕获的是那个脚本运行时 TRIPS_RAW 的状态,而模块 04 是在那份快照之上构建 dbt 流水线的。 自迁移脚本最后一批之后写入 Snowflake 的每一条行程都只存在于 Snowflake 中,ClickHouse 缺了尾巴。步骤 3 的 --resume 追赶那一遍 恰好弥合这个间隔:它读取 ClickHouse 中已有的 max(pickup_at),只拉取 在那之后写入的行,所以弥合一个自模块 01 起就存在的间隔只需数秒 到数分钟,而不是原始批量迁移所花的 40-50 分钟。跳过它就直接切换, ClickHouse 会永久丢掉落在那个间隔里的所有行程,这是步骤 4 专门要抓住的一种 无声的一致性故障,但前提是步骤 3 按顺序跑过了。

agg_hourly_zone_trips 在这里、也只在这里被填充。 它自模块 03 起一直是空的, 这是设计如此:它的增量过滤条件是 WHERE pickup_at >= now() - INTERVAL 2 HOUR, 只会匹配实时生产者写入的行,而在本模块之前,唯一在写的生产者 是 Snowflake 那一侧的。一旦步骤 3 启动 ClickHouse 生产者,新的数据行终于落在 那个两小时窗口内,于是这张表,以及每一张由它支撑的仪表板图表, 在本实验中第一次不再空白。

步骤 1:添加 ClickHouse 仪表板

如果你想走完整的手动搭建流程,在 Superset 界面里一步步创建全部 7 个数据集、18 张图表和 4 个仪表板,请照做 ClickHouse 上的 Superset。若想跳过 手动步骤、一次性把所有东西导入:

cd "$(git rev-parse --show-toplevel)/workshop_public/snowflake_migration_lab/03-migrate-to-clickhouse"
source .env && source .clickhouse_state
bash superset/add_clickhouse_connection.sh

已提交的 superset/dashboards/dashboard_export_*.zip 中,ClickHouse 主机名被 脱敏为 your-instance.clickhouse.cloud。 上面的快捷路径会在导入前用 .env 中的值修补该 URI,所以这一步是透明的,你不会察觉到脱敏。但 如果你改为通过 Superset 界面手动导入 ZIP,它创建的数据库连接 将连不上,之后你必须编辑那个连接,指向你 真实的 CLICKHOUSE_HOST 和凭据。具体做法见 ClickHouse 上的 Superset。

验证:

打开 http://localhost:8088(admin / admin)。在 Dashboards 下,你应该看到共 7 个,3 个 Snowflake 仪表板和 4 个以 CH — 为前缀的仪表板。

步骤 2:运行基准测试

对 Snowflake 和 ClickHouse 背靠背运行全部七个查询,并比较 实际耗时:

cd "$(git rev-parse --show-toplevel)/workshop_public/snowflake_migration_lab/03-migrate-to-clickhouse"
source .env && source .clickhouse_state
./scripts/run_benchmark.sh

预期输出:

━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━
  NYC Taxi Lab — Query Benchmark: Snowflake vs ClickHouse
  (median of 3 runs each)
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━
Query                                   Snowflake     ClickHouse    Speedup
────────────────────────────────────────────────────────────────────────
Q1  Hourly revenue by borough           5.0s          0.7s          6x
Q2  Rolling 7-day avg distance          5.5s          0.8s          6x
Q3  Top 10 trips (QUALIFY→subquery)     5.0s          0.7s          6x
Q4  Driver ratings (JSON flatten)       5.4s          0.8s          6x
Q5  Surge pricing (VARIANT)             5.1s          0.7s          6x
Q6  Hourly aggregation (MERGE→RMT)      5.9s          0.8s          7x
Q7  CDC/live data freshness             7.9s          0.8s          9x
────────────────────────────────────────────────────────────────────────
Total                                   40.1s         5.6s          7x avg
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

这些数字来自一次有代表性的运行,不是保证,你自己的数字会随 warehouse 规模、 ClickHouse Cloud 层级,以及当时还有什么别的负载在跑而 变化。

脚本会把每次运行写入 workshop_public/snowflake_migration_lab/03-migrate-to-clickhouse/scripts/benchmark_results_<timestamp>.csv。 把它留在磁盘上,它是模块 06 考核所需的两个文件之一,所以在步骤 5 拆除资源时 不要删掉它。

步骤 3:切换到 ClickHouse

迁移间隔。 Snowflake 生产者在整个实验期间一直在运行,以约每分钟 60 条行程的 速度写入。模块 03 的迁移脚本捕获的是那个脚本运行时 TRIPS_RAW 的状态,此后写入的数据行只存在于 Snowflake 中。在切换写入路径之前 先弥合那个间隔。

按顺序执行下面全部四步,顺序正是保证两套系统一致的关键:

cd "$(git rev-parse --show-toplevel)/workshop_public/snowflake_migration_lab/03-migrate-to-clickhouse"
source .env && source .clickhouse_state
source .venv/bin/activate

# Step 1: Stop the Snowflake producer (freeze the dataset)
docker stop nyc_taxi_producer

# Step 2: Catch up the delta — only migrates rows with pickup_at newer than
# what's already in ClickHouse. Runs in seconds to minutes, not the original
# 40-50 minutes, because only the gap rows move.
python scripts/02_migrate_trips.py --resume

# Step 3: Refresh the analytics tables with the newly migrated rows
cd "$(git rev-parse --show-toplevel)/workshop_public/snowflake_migration_lab/03-migrate-to-clickhouse/dbt/nyc_taxi_dbt_ch"
dbt run
cd "$(git rev-parse --show-toplevel)/workshop_public/snowflake_migration_lab/03-migrate-to-clickhouse"

# Step 4: Start the ClickHouse producer
source .env && source .clickhouse_state
./scripts/03_cutover.sh

不要跳过 Step 2。 --resume 会读取 ClickHouse 中已有的 max(pickup_at),并给 Snowflake 查询加上一个 WHERE PICKUP_DATETIME > <watermark> 过滤条件,因此它只传输 在模块 03 原始迁移运行期间及之后写入的那些行,也就是 ClickHouse 从未见过的那些。不做这一步就切换,ClickHouse 会永久丢掉 落在那个窗口里的所有行程。步骤 4 的一致性检查正是为抓住这一点而设计的, 但前提是这一步先跑过。

./scripts/03_cutover.sh 会提示 Type "cutover" to confirm,然后把第 1-3 步作为 它自己的安全网重做一遍:它再次停掉 Snowflake 生产者(如果你已经停过,这是空操作)、 再跑一次 dbt run,然后构建并启动 ClickHouse 生产者 (nyc_taxi_ch_producer)。生产者启动后三十秒,它会确认新的数据行 正在落入 default.trips_raw,并再跑一次 dbt,正是这次运行终于给了 agg_hourly_zone_trips 它的第一批数据行,弥合了模块 04 有意留下的那个空缺。

验证:

-- Most recent trip should be within the last 60 seconds
SELECT max(pickup_at) AS most_recent_trip FROM default.trips_raw;

-- Row count should be increasing — wait 60 seconds and run again
SELECT count() FROM default.trips_raw;

-- agg_hourly_zone_trips should now have rows for the first time in the lab
SELECT count() FROM analytics.agg_hourly_zone_trips;
docker ps | grep nyc_taxi_ch_producer   # should show running

保持分析层新鲜。 fact_trips 和 agg_hourly_zone_trips 是 dbt 增量模型,它们不会自动刷新。03_cutover.sh 在确认生产者已上线后会跑一次 dbt run,但随着新行程不断积累,仪表板会逐渐变旧; 想要最新数字时,就从 dbt/nyc_taxi_dbt_ch 再跑一次 dbt run(在 生产环境中你会把它排期,cron、Airflow、dbt Cloud,但在实验里按需运行就 够了)。相比之下,analytics.mv_live_trip_feed 是一个可刷新的 materialized view,模块 04 的 dbt run 已经用 engine = 'ReplacingMergeTree(refreshed_at)' 把它建好了,但本实验从不开启它的刷新间隔: 那条能让它自行重新执行的 MODIFY REFRESH EVERY 30 SECOND 语句 只以注释形式存在于模型文件中。启用它只是一条你自己去执行的 ALTER TABLE 语句;不启用的话,mv_live_trip_feed 只在 dbt 构建它的那一次更新。

反向切换,如果你需要撤销这一步、回到 Snowflake 生产者:

docker stop nyc_taxi_ch_producer
cd "$(git rev-parse --show-toplevel)/workshop_public/snowflake_migration_lab/01-setup-snowflake/superset"
docker-compose --env-file ../.env up -d producer

步骤 4:验证一致性

既然 --resume 追赶那一遍已经跑过、ClickHouse 生产者也已激活,两套 系统现在应该是一致的。来确认一下:

cd "$(git rev-parse --show-toplevel)/workshop_public/snowflake_migration_lab/03-migrate-to-clickhouse"
source .env && source .clickhouse_state
bash scripts/01_verify_migration.sh

预期输出:

━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━
  Migration Parity Check
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

  ✓ ClickHouse default.trips_raw: 50,008,250 rows

  ✓ Snowflake NYC_TAXI_DB.RAW.TRIPS_RAW: 50,008,250 rows
  ✓ Row count parity: PASS  (difference: 0 rows = 0.0000%)

  ✓ trip_metadata populated: 50,008,250 non-empty rows
  pickup_at range: 2022-03-30   2026-03-31

  ✓ ClickHouse has 50,008,250 rows — migration looks complete
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

你自己的行数会不一样;要紧的是那一行一致性结论。到这一步 Snowflake 生产者已经 停了,所以那边不再有新行落入,行数应当完全吻合,或者在 --resume 那一遍期间恰有一批数据在途时相差寥寥几行,远远处于脚本所检查的 0.01% 阈值之内。

如果一致性检查失败(差异大于 0.01%),说明间隔没有完全 弥合,再跑一次追赶并重新检查:

cd "$(git rev-parse --show-toplevel)/workshop_public/snowflake_migration_lab/03-migrate-to-clickhouse"
source .env && source .clickhouse_state
python scripts/02_migrate_trips.py --resume
bash scripts/01_verify_migration.sh

步骤 5:拆除资源

在拆除任何东西之前,先确认迁移处于正确的最终状态:

检查项命令预期
行数一致性bash scripts/01_verify_migration.sh行数吻合度 ≥ 99.9%
dbt 测试dbt test(在 dbt/nyc_taxi_dbt_ch 目录下)所有测试通过
Superset 仪表板打开 http://localhost:8088可见 7 个仪表板(3 个 SF + 4 个 CH)
基准测试结果cat scripts/benchmark_results_<timestamp>.csv全部 7 个查询都有提速倍数
cd "$(git rev-parse --show-toplevel)/workshop_public/snowflake_migration_lab/03-migrate-to-clickhouse"
source .env && source .clickhouse_state
bash scripts/01_verify_migration.sh

cd "$(git rev-parse --show-toplevel)/workshop_public/snowflake_migration_lab/03-migrate-to-clickhouse/dbt/nyc_taxi_dbt_ch"
dbt test

四项全部通过后,有两个文件就是模块 06 所需的一切,而且它们都能在拆除后留存: workshop_public/snowflake_migration_lab/02-plan-and-design/migration-plan.md 和 workshop_public/snowflake_migration_lab/03-migrate-to-clickhouse/scripts/benchmark_results_<timestamp>.csv。 模块 06 是一场开卷的纸面考核,约 60 分钟,除此之外什么都不 需要,没有任何理由为了一场笔试而让一个付费的 ClickHouse Cloud 服务一直 开着。 把这两个文件的内容复制或记录到你能拿到的地方, 然后拆除所有资源:

cd "$(git rev-parse --show-toplevel)/workshop_public/snowflake_migration_lab/03-migrate-to-clickhouse"
source .env && ./teardown.sh

这会销毁 ClickHouse Cloud 服务(通过 terraform destroy),以及在执行过切换的情况下 销毁 ClickHouse 行程生产者容器。

第 1 部分的 Snowflake 资源不会被这个脚本拆除。 请单独拆除 Snowflake 一侧:

cd "$(git rev-parse --show-toplevel)/workshop_public/snowflake_migration_lab/01-setup-snowflake"
source .env && ./teardown.sh

如何确认已完成

到本模块的这个位置,你应当已按顺序确认过:

  • 一致性检查通过:步骤 4 的 01_verify_migration.sh 报告 PASS,且 行数差异低于 0.01%。
  • 基准测试 CSV 在磁盘上:步骤 2 写出了 benchmark_results_<timestamp>.csv,其中全部 7 个查询都带有提速倍数,而且你在步骤 5 拆除之前把它保留了下来。
  • 7 个仪表板都在:步骤 1 的 Superset 检查显示 3 个 Snowflake 仪表板和 4 个 CH — 仪表板并列在一起。
  • ClickHouse 生产者在写入:步骤 3 的验证代码块显示 default.trips_raw 行数在增长、nyc_taxi_ch_producer 在运行,然后步骤 5 才为拆除而停掉它。

如果其中任何一项当时没有成立,请回到对应的步骤,而不是现在 重跑这些检查,步骤 5 已经销毁了 ClickHouse Cloud 服务, 如果做过切换,连生产者容器也一起销毁了。

结束状态

迁移已完成并已量化:5000 万行从 Snowflake 搬到了 ClickHouse 并验证了一致性,七个查询做了正面对比基准测试且 ClickHouse 在每一个上都更快,BI 层已重建,4 个 ClickHouse 仪表板与原有的 3 个 Snowflake 仪表板并列,写入路径也已从 Snowflake 永久切换到 ClickHouse。 两个云环境都已拆除,没有 ClickHouse Cloud 服务、没有 ClickHouse 生产者容器,而且在第 1 部分的拆除也执行之后,也没有 Snowflake warehouse 了。

有两个文件在拆除后留存,并且是模块 06 所需的一切: workshop_public/snowflake_migration_lab/02-plan-and-design/migration-plan.md 和 workshop_public/snowflake_migration_lab/03-migrate-to-clickhouse/scripts/benchmark_results_<timestamp>.csv。 模块 06 是一场开卷书面考核,带上这两个文件,别的都不用。

本页内容

Track your progress?

Optional. We email a link to confirm your address; progress records once you open it.

Please use your work email address, not a personal one.

Progress tracking also requires accepting the current Terms of Service in Privacy settings.

ZH