Snowflake MigrationClickHouse Workshops

示例解答:一份已完成的计划

针对 NYC taxi 负载填写完整的迁移计划,供你写完自己的版本后对照。

这是把五份练习册 (1、 2、 3、 4、 5) 全部应用到 NYC Taxi 负载后的完整解答。你可以用它来:

  • 在完成每一节后核对自己的练习册答案
  • 理解第 3 部分所实现的那些决策背后的推理
  • 如果你的选择不同,用它与第 3 部分的 Decision Alignment 表格做对照

这是答案册,不要把它当作你自己的计划来填写。请填写 migration-plan.md。


完成情况清单

  • 引擎选择:已完成
  • 排序键设计:已完成
  • Schema 转换:已完成
  • 迁移批次计划:已完成
  • dbt 模型设计:已完成

第 1 节:画像概览

指标值
表总数7(TRIPS_RAW、FACT_TRIPS、AGG_HOURLY_ZONE_TRIPS、DIM_TAXI_ZONES、DIM_PAYMENT_TYPE、DIM_VENDOR、DIM_DATE)
视图总数2(STG_TRIPS、STG_TAXI_ZONES)
Streams1(TRIPS_RAW 上的 TRIPS_CDC_STREAM)
Tasks2(CDC_CONSUME_TASK、HOURLY_AGG_TASK)
TRIPS_RAW 总行数~50,000,000
日期范围以搭建时间为终点的 4 年滚动窗口
VARIANT 列1(TRIPS_RAW.TRIP_METADATA)
检出的 QUALIFY 使用处1(Q3 查询)
检出的 MERGE INTO 使用处2(HOURLY_AGG_TASK、CDC_CONSUME_TASK)

第 2 节:对象清单

对象类型Schema行数复杂度等级说明
trips_raw表raw~50MB以 _synced_at 作为版本列的 RMT;批量加载与 CDC 存在重叠,需要去重;stg_trips 必须使用 FINAL
stg_tripsdbt 视图staging,B对 TRIP_METADATA 使用 JSONExtract;需要测试 JSON 路径
stg_taxi_zonesdbt 视图staging,A直通;很简单
int_trips_enricheddbt Ephemeralstaging,ACTE;SQL 差异在父模型中处理
fact_tripsdbt Incrementalanalytics~50MCRMT 引擎;delete_insert;需改写 QUALIFY;必须用 FINAL
agg_hourly_zone_tripsdbt Incrementalanalytics~140KBRMT;滚动 2 小时重算窗口;需仔细测试分区边界
dim_taxi_zonesdbt 表analytics265A静态参照数据;完全重载;很简单
dim_payment_typedbt 表analytics6A静态参照数据;很简单
dim_vendordbt 表analytics3A静态参照数据;很简单
taxi_zones_dict字典analytics265BClickHouse 专属语法;查询时使用 dictGet()
mv_hourly_revenue可刷新 MVanalytics,BREFRESH EVERY 语法;需验证原子替换
TRIPS_CDC_STREAM / CDC_CONSUME_TASKSnowflake Stream + Task,,DClickHouse 无对应物;在第 3 部分改由生产者直接切流替代

第 3 节:引擎选择决策

表引擎版本列理由
trips_rawReplacingMergeTree(_synced_at)_synced_atPython 迁移脚本(scripts/02_migrate_trips.py)可能重试某一批次,从而重复插入同一个 trip_id。切流之后,线上生产者也可能在瞬时故障后重试。_synced_at DateTime DEFAULT now() 会在 INSERT 时自动赋值,后来的重试会有更大的时间戳,因此 RMT 会保留最新那次写入。stg_trips 使用 FINAL 查询,以在任何下游模型运行之前完成去重。
fact_tripsReplacingMergeTree(updated_at)updated_at行程可能被修正(车费调整、状态变更)。同一个 trip_id 会带着更新后的值再次插入。updated_at 在每次修正时单调递增,RMT 去重时取值更大的一行胜出。查询时始终使用 FINAL。
agg_hourly_zone_tripsReplacingMergeTree(updated_at)updated_atdbt 会重算最近 2 小时并重新插入。没有 RMT 的话,旧聚合与新聚合会累积并导致重复计数。每次 dbt run 时把 updated_at 设为 now(),可确保最新值胜出。
dim_taxi_zonesMergeTree(),由 dbt 完全重载(原子表交换(完全重建))。不可能累积重复。无需去重。
dim_payment_typeMergeTree(),同上,完全重载。
dim_vendorMergeTree(),同上,完全重载。
mv_hourly_revenueMergeTree(),REFRESHABLE MV 在每次 REFRESH 时原子地替换整个结果集。没有 upsert。

第 4 节:排序键设计

表ORDER BY理由
trips_raw(pickup_at, trip_id)时间范围扫描首先按 pickup_at 过滤。trip_id 是 RMT 的去重键,它必须出现在 ORDER BY 中,RMT 才能识别哪些行是重复的。pickup_at 放在前面,因为分析型范围扫描占主导;trip_id 放在最后,因为它基数很高,只承担唯一性判别的作用。
fact_trips(toStartOfMonth(pickup_at), pickup_at, trip_id)全部 7 个分析查询都按 pickup_at 过滤。月份前缀把同一自然月的数据聚到相邻数据块中,让按月聚合能做粗粒度的块跳过,而无需再加 PARTITION BY。trip_id 放最后,既满足 RMT 的唯一性要求,又不破坏块跳过。
agg_hourly_zone_trips(hour_bucket, zone_id)Q6(以及所有聚合查询)都按 hour_bucket 和 zone_id 过滤。hour_bucket 有约 3.5 万个不同取值;zone_id 有 265 个。hour_bucket 放前面,因为时间范围扫描是主要访问模式。zone_id 放第二位,用于二级过滤。
dim_taxi_zones(location_id)265 行 = 一个 granule。ORDER BY 对性能没有影响。用 location_id 作为 join 键符合惯例,也便于阅读。

第 5 节:Schema 转换说明

列Snowflake 类型ClickHouse 类型决策依据
TRIP_METADATAVARIANTString原样保留原始 JSON。JSONExtract* 可在查询时处理任意路径。Map(String,String) 会丢失嵌套结构;Tuple 需要固定 schema。对任意 JSON 来说,String 是稳妥的选择。
PICKUP_DATETIME / PICKUP_ATTIMESTAMP_NTZ(9)DateTime64(3, 'UTC')对行程时间戳来说毫秒精度已足够。纳秒(9)过度了。'UTC' 让时区显式化,避免时间范围聚合中出现夏令时相关的意外。
PICKUP_LOCATION_IDINTEGERUInt16取值范围 1–265。UInt8 最大值为 255(太小)。UInt16 最大值为 65535(合适)。相比 Int32 的 4 字节只占 2 字节,在 5000 万行上每列可省下约 95MB 未压缩空间。
VENDOR_IDINTEGERUInt8取值范围 1–3。UInt8 最大值 255,合适。每行 1 字节。
DRIVER_RATINGFLOATNullable(Float32)经常为 NULL(并非所有行程都有评分)。Nullable 保留了正确的 null 语义。Float32 对 1.0–5.0 的范围已足够。Float64 只会浪费存储,并不带来有意义的精度提升。
UPDATED_ATTIMESTAMP_NTZ(9)DateTime64(3, 'UTC')ReplacingMergeTree 的版本列。必须用 DateTime64 而不是 DateTime,同一秒内的两次修正在秒级精度下会变得不确定。毫秒精度保证去重顺序正确。

需要做的函数转换

Snowflake 表达式ClickHouse 等价写法
DATE_TRUNC('hour', pickup_at)toStartOfHour(pickup_at)
DATEADD('day', -7, CURRENT_DATE)today() - 7
DATEDIFF('minute', pickup_at, dropoff_at)dateDiff('minute', pickup_at, dropoff_at)
TRIP_METADATA:driver.rating::FLOATJSONExtractFloat(trip_metadata, 'driver', 'rating')
TRIP_METADATA:surge_multiplier::FLOATJSONExtractFloat(trip_metadata, 'surge_multiplier')
QUALIFY ROW_NUMBER() OVER (PARTITION BY pickup_location_id ORDER BY fare_amount DESC) <= 10SELECT ... FROM (SELECT ..., ROW_NUMBER() OVER (...) AS rn FROM ...) WHERE rn <= 10
MERGE INTO fact_trips ... WHEN MATCHED THEN UPDATEdbt 的 delete_insert 增量模式,先删除键匹配的行,再插入所有新行

第 6 节:迁移批次

批次对象依赖说明
批次 0trips_raw(schema)、dim_taxi_zones、dim_payment_type、dim_vendor无dbt 创建空表。维度表立即用静态参照数据填充(不依赖 trips)。运行:dbt run --select trips_raw dim_*
批次 1Python 批量加载(scripts/02_migrate_trips.py)批次 0(trips_raw 的 schema 必须已存在)从 Snowflake TRIPS_RAW 迁移 5000 万行。可用 --resume 断点续传。用 scripts/01_verify_migration.sh 核对行数。约 40–50 分钟。
批次 2stg_trips、stg_taxi_zones、int_trips_enriched、fact_trips、agg_hourly_zone_trips批次 1 已完成(trips_raw 已填充)+ 批次 0(维度表已存在)完整的 dbt run。stg_trips 读取 trips_raw;int_trips_enriched 与维度表关联;fact_trips 和 agg_hourly_zone_trips 在其之上构建。
批次 3taxi_zones_dict、mv_live_trip_feed批次 2(字典需要 dim_taxi_zones 已填充;MV 需要 fact_trips 已填充)字典通过 scripts/04_create_dictionary.sql 创建。可刷新 MV 通过 dbt 模型创建;之后启用其刷新间隔是一次手动的 ALTER TABLE ... MODIFY REFRESH 操作,dbt 不会自动执行。
批次 4生产者切流(scripts/03_cutover.sh)批次 1 已完成(批量加载已核验)+ 批次 2 已完成(analytics 层已构建)停止 Snowflake 生产者;启动直接写入 ClickHouse Cloud 的 ClickHouse 生产者;运行 dbt run,用实时数据填充 agg_hourly_zone_trips。

风险登记表(等级 C/D 的对象)

对象风险验证方法
fact_trips在合并滞后期间,不带 FINAL 的查询会重复计数。delete_insert 的分区范围必须限定在 ORDER BY 前缀内,以免删除非目标分区。SELECT COUNT(*) FINAL 与 Snowflake 一致(允许 CDC 滞后带来的差异)。运行 dbt test。对比两套系统的 Q3 结果。检查是否有重复 trip_id:SELECT trip_id, count() FROM fact_trips GROUP BY trip_id HAVING count() > 1 LIMIT 10。
agg_hourly_zone_trips滚动 2 小时重算窗口必须正确界定删除范围。范围过宽会删掉旧聚合;过窄则会残留过期聚合。抽查特定的 (hour_bucket, zone_id) 组合并与 Snowflake 对比。核对同一时段内所有 zone 的 trip_count 总数是否与 Snowflake 的 AGG_HOURLY_ZONE_TRIPS 一致。
生产者切流迁移脚本中途被打断会留下行数缺口;用 --resume 重跑补齐。切流后生产者的重试可能重复插入已在 ClickHouse 中的行程。scripts/01_verify_migration.sh,检查 Snowflake 与 ClickHouse 的行数是否一致。ReplacingMergeTree(_synced_at) 让重复插入具有幂等性。

第 7 节:已知的方言差异

  • QUALIFY,影响:Q3(queries/q03_top_trips_qualify.sql)
  • VARIANT 冒号路径,影响:Q4、Q5(TRIP_METADATA 的 JSON 访问)
  • LATERAL FLATTEN,本负载未使用;VARIANT 通过冒号路径访问,而非 FLATTEN
  • MERGE INTO,影响:dbt 增量模型(fact_trips、agg_hourly_zone_trips)
  • Snowflake Streams → 生产者切流(切流后实时写入直接进入 ClickHouse)
  • 日期函数差异,影响:Q1(DATE_TRUNC)、Q3(DATEADD)、Q4(DATEDIFF)

第 8 节:迁移策略

数据搬迁: Python 迁移脚本(scripts/02_migrate_trips.py)

为什么选 Python 脚本,而不是对象存储中转或 ClickPipes?

  • remoteSecure() 用于 ClickHouse 到 ClickHouse 的数据传输,在这里不适用。
  • 对象存储中转(Snowflake → S3 → ClickHouse 的 S3 表函数)可行,但会增加复杂度:需要开通 S3 存储桶、IAM 角色以及 Snowflake 的 COPY INTO,对一个实验课来说是不必要的开销。
  • ClickPipes 不支持以 Snowflake 作为数据源。它支持的数据源是 Kafka、S3、Kinesis、PostgreSQL CDC 和 MySQL CDC。
  • 这个 Python 脚本使用 snowflake-connector-python 和 clickhouse-connect,都是实验课已经安装的包。它会实时显示进度,支持 --resume 应对中断,而且代码完全可供查阅。

增量策略(dbt): delete_insert

为什么选 delete_insert,而不是 append 或 merge 策略?

  • append 插入新行而不触碰已有行。对于行会被更新的 fact_trips 来说,这会产生重复。不正确。
  • merge(如果可用)最接近 Snowflake 的 MERGE INTO,但 dbt-clickhouse 的 merge 策略在与 ReplacingMergeTree 配合时有局限,并非推荐做法。
  • delete_insert 会删除落在本批次键范围内的行,然后插入所有新行。它是幂等的(重跑得到相同结果),既能处理插入也能处理更新,并且与 ReplacingMergeTree 配合正确。它是 dbt-clickhouse 社区对 upsert 模式的标准推荐。

第 9 节:切流标准

标准阈值度量方式
行数一致性匹配度 ≥ 99.9%(切流后 CH ≥ SF 属预期)scripts/01_verify_migration.sh
校验和一致性1 万行样本上 MD5 一致scripts/02_validate_parity.sql
dbt 测试通过率100%在 dbt/nyc_taxi_dbt_ch 中执行 dbt test
查询结果一致性全部 7 个查询返回相同结果(在浮点误差范围内)在 scripts/run_benchmark.sh 的输出中手工对比


第 10 节:dbt 模型设计

物化方式的选择

模型物化方式原因
stg_tripsview读取并清洗 trips_raw;该模型不会被更新;零存储成本;始终反映当前源状态
stg_taxi_zonesview同上,对一张源表做直通式清洗
int_trips_enrichedephemeral纯 join 逻辑,只被 fact_trips 使用;内联为 CTE 可避免多出一张冗余物理表;没有模型直接查询它
fact_tripsincremental行程可能事后被修正;每次运行只应处理新增和更新的行
agg_hourly_zone_tripsincremental滚动 2 小时重算本身就是增量模式,只处理近期的行,而不是全部 5000 万行
dim_taxi_zonestable265 个静态 zone;每次 dbt run 都通过原子表交换(完全重建)整体重建;不需要局部更新
dim_payment_typetable6 种静态类型;理由同 dim_taxi_zones
dim_vendortable3 个 vendor;理由相同

引擎配置

模型ENGINE版本列原因
fact_tripsReplacingMergeTree(updated_at)updated_at行程可能被修正;每次插入时把 updated_at 设为 now(),意味着 RMT 后台去重时最新版本胜出;delete_insert 是主要的正确性路径,RMT 是安全网
agg_hourly_zone_tripsReplacingMergeTree(updated_at)updated_at滚动重算会为相同的 (hour_bucket, zone_id) 组合重新插入聚合;RMT 确保过期聚合在后台合并时被移除
dim_taxi_zonesMergeTree(),由 dbt 完全重载意味着每次运行都做原子表交换(完全重建);不可能累积重复;无需去重
dim_payment_typeMergeTree(),与 dim_taxi_zones 相同
dim_vendorMergeTree(),与 dim_taxi_zones 相同

增量策略

模型unique_keyincremental_strategy增量过滤条件为什么用这个过滤条件?
fact_tripstrip_iddelete_insertWHERE updated_at > (SELECT max(updated_at) FROM {{ this }})以 updated_at 作为高水位线,既能捕获新增行程,也能捕获被修正的行程(车费调整会以相同的 trip_id、相同的 pickup_at、但更新的 updated_at 重新插入);若用 pickup_at 做水位线,会静默漏掉这些修正
agg_hourly_zone_trips[hour_bucket, zone_id]delete_insertWHERE pickup_at >= now() - INTERVAL 2 HOUR滚动 2 小时窗口强制对边界小时重新聚合,使不完整小时的计数总能被修正;若用 max(pickup_at) 高水位线,边界小时会被永久低估

FINAL 的放置

模型FROM 子句中是否加 FINAL?原因
stg_trips是,FROM trips_raw FINALtrips_raw 是 ReplacingMergeTree;由于迁移脚本重试或切流后生产者重试,它可能出现重复的 trip_id 行。stg_trips 是唯一的执行点:在这里去重,让每个下游模型(int_trips_enriched、fact_trips、agg_hourly_zone_trips)都拿到干净数据
int_trips_enriched否它读取的是 stg_trips(一个视图),而不是 RMT 表;FINAL 对视图无意义
fact_trips否(在模型体内)每次运行完成后,delete_insert 都会让 fact_trips 保持干净;在模型内部加 FINAL 会被浪费地施加到那个从 {{ this }} 读取 max(updated_at) 的 is_incremental() 子查询上。看板和 dbt 测试在直接查询 fact_trips 时会在外部使用 FINAL

这是完成后的示例。你的 migration-plan.md 应当与这里的关键决策一致,或者明确记录你为什么做了不同的选择。

本页内容

ZH