示例解答:一份已完成的计划
针对 NYC taxi 负载填写完整的迁移计划,供你写完自己的版本后对照。
这是把五份练习册
(1、
2、
3、
4、
5)
全部应用到 NYC Taxi 负载后的完整解答。你可以用它来:
- 在完成每一节后核对自己的练习册答案
- 理解第 3 部分所实现的那些决策背后的推理
- 如果你的选择不同,用它与第 3 部分的 Decision Alignment 表格做对照
这是答案册,不要把它当作你自己的计划来填写。请填写 migration-plan.md。
| 指标 | 值 |
|---|
| 表总数 | 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) |
| Streams | 1(TRIPS_RAW 上的 TRIPS_CDC_STREAM) |
| Tasks | 2(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) |
| 对象 | 类型 | Schema | 行数 | 复杂度等级 | 说明 |
|---|
trips_raw | 表 | raw | ~50M | B | 以 _synced_at 作为版本列的 RMT;批量加载与 CDC 存在重叠,需要去重;stg_trips 必须使用 FINAL |
stg_trips | dbt 视图 | staging | , | B | 对 TRIP_METADATA 使用 JSONExtract;需要测试 JSON 路径 |
stg_taxi_zones | dbt 视图 | staging | , | A | 直通;很简单 |
int_trips_enriched | dbt Ephemeral | staging | , | A | CTE;SQL 差异在父模型中处理 |
fact_trips | dbt Incremental | analytics | ~50M | C | RMT 引擎;delete_insert;需改写 QUALIFY;必须用 FINAL |
agg_hourly_zone_trips | dbt Incremental | analytics | ~140K | B | RMT;滚动 2 小时重算窗口;需仔细测试分区边界 |
dim_taxi_zones | dbt 表 | analytics | 265 | A | 静态参照数据;完全重载;很简单 |
dim_payment_type | dbt 表 | analytics | 6 | A | 静态参照数据;很简单 |
dim_vendor | dbt 表 | analytics | 3 | A | 静态参照数据;很简单 |
taxi_zones_dict | 字典 | analytics | 265 | B | ClickHouse 专属语法;查询时使用 dictGet() |
mv_hourly_revenue | 可刷新 MV | analytics | , | B | REFRESH EVERY 语法;需验证原子替换 |
TRIPS_CDC_STREAM / CDC_CONSUME_TASK | Snowflake Stream + Task | , | , | D | ClickHouse 无对应物;在第 3 部分改由生产者直接切流替代 |
| 表 | 引擎 | 版本列 | 理由 |
|---|
trips_raw | ReplacingMergeTree(_synced_at) | _synced_at | Python 迁移脚本(scripts/02_migrate_trips.py)可能重试某一批次,从而重复插入同一个 trip_id。切流之后,线上生产者也可能在瞬时故障后重试。_synced_at DateTime DEFAULT now() 会在 INSERT 时自动赋值,后来的重试会有更大的时间戳,因此 RMT 会保留最新那次写入。stg_trips 使用 FINAL 查询,以在任何下游模型运行之前完成去重。 |
fact_trips | ReplacingMergeTree(updated_at) | updated_at | 行程可能被修正(车费调整、状态变更)。同一个 trip_id 会带着更新后的值再次插入。updated_at 在每次修正时单调递增,RMT 去重时取值更大的一行胜出。查询时始终使用 FINAL。 |
agg_hourly_zone_trips | ReplacingMergeTree(updated_at) | updated_at | dbt 会重算最近 2 小时并重新插入。没有 RMT 的话,旧聚合与新聚合会累积并导致重复计数。每次 dbt run 时把 updated_at 设为 now(),可确保最新值胜出。 |
dim_taxi_zones | MergeTree() | , | 由 dbt 完全重载(原子表交换(完全重建))。不可能累积重复。无需去重。 |
dim_payment_type | MergeTree() | , | 同上,完全重载。 |
dim_vendor | MergeTree() | , | 同上,完全重载。 |
mv_hourly_revenue | MergeTree() | , | REFRESHABLE MV 在每次 REFRESH 时原子地替换整个结果集。没有 upsert。 |
| 表 | 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 键符合惯例,也便于阅读。 |
| 列 | Snowflake 类型 | ClickHouse 类型 | 决策依据 |
|---|
TRIP_METADATA | VARIANT | String | 原样保留原始 JSON。JSONExtract* 可在查询时处理任意路径。Map(String,String) 会丢失嵌套结构;Tuple 需要固定 schema。对任意 JSON 来说,String 是稳妥的选择。 |
PICKUP_DATETIME / PICKUP_AT | TIMESTAMP_NTZ(9) | DateTime64(3, 'UTC') | 对行程时间戳来说毫秒精度已足够。纳秒(9)过度了。'UTC' 让时区显式化,避免时间范围聚合中出现夏令时相关的意外。 |
PICKUP_LOCATION_ID | INTEGER | UInt16 | 取值范围 1–265。UInt8 最大值为 255(太小)。UInt16 最大值为 65535(合适)。相比 Int32 的 4 字节只占 2 字节,在 5000 万行上每列可省下约 95MB 未压缩空间。 |
VENDOR_ID | INTEGER | UInt8 | 取值范围 1–3。UInt8 最大值 255,合适。每行 1 字节。 |
DRIVER_RATING | FLOAT | Nullable(Float32) | 经常为 NULL(并非所有行程都有评分)。Nullable 保留了正确的 null 语义。Float32 对 1.0–5.0 的范围已足够。Float64 只会浪费存储,并不带来有意义的精度提升。 |
UPDATED_AT | TIMESTAMP_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::FLOAT | JSONExtractFloat(trip_metadata, 'driver', 'rating') |
TRIP_METADATA:surge_multiplier::FLOAT | JSONExtractFloat(trip_metadata, 'surge_multiplier') |
QUALIFY ROW_NUMBER() OVER (PARTITION BY pickup_location_id ORDER BY fare_amount DESC) <= 10 | SELECT ... FROM (SELECT ..., ROW_NUMBER() OVER (...) AS rn FROM ...) WHERE rn <= 10 |
MERGE INTO fact_trips ... WHEN MATCHED THEN UPDATE | dbt 的 delete_insert 增量模式,先删除键匹配的行,再插入所有新行 |
| 批次 | 对象 | 依赖 | 说明 |
|---|
| 批次 0 | trips_raw(schema)、dim_taxi_zones、dim_payment_type、dim_vendor | 无 | dbt 创建空表。维度表立即用静态参照数据填充(不依赖 trips)。运行:dbt run --select trips_raw dim_* |
| 批次 1 | Python 批量加载(scripts/02_migrate_trips.py) | 批次 0(trips_raw 的 schema 必须已存在) | 从 Snowflake TRIPS_RAW 迁移 5000 万行。可用 --resume 断点续传。用 scripts/01_verify_migration.sh 核对行数。约 40–50 分钟。 |
| 批次 2 | stg_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 在其之上构建。 |
| 批次 3 | taxi_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。 |
| 对象 | 风险 | 验证方法 |
|---|
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) 让重复插入具有幂等性。 |
数据搬迁: 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 模式的标准推荐。
| 标准 | 阈值 | 度量方式 |
|---|
| 行数一致性 | 匹配度 ≥ 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 的输出中手工对比 |
| 模型 | 物化方式 | 原因 |
|---|
stg_trips | view | 读取并清洗 trips_raw;该模型不会被更新;零存储成本;始终反映当前源状态 |
stg_taxi_zones | view | 同上,对一张源表做直通式清洗 |
int_trips_enriched | ephemeral | 纯 join 逻辑,只被 fact_trips 使用;内联为 CTE 可避免多出一张冗余物理表;没有模型直接查询它 |
fact_trips | incremental | 行程可能事后被修正;每次运行只应处理新增和更新的行 |
agg_hourly_zone_trips | incremental | 滚动 2 小时重算本身就是增量模式,只处理近期的行,而不是全部 5000 万行 |
dim_taxi_zones | table | 265 个静态 zone;每次 dbt run 都通过原子表交换(完全重建)整体重建;不需要局部更新 |
dim_payment_type | table | 6 种静态类型;理由同 dim_taxi_zones |
dim_vendor | table | 3 个 vendor;理由相同 |
| 模型 | ENGINE | 版本列 | 原因 |
|---|
fact_trips | ReplacingMergeTree(updated_at) | updated_at | 行程可能被修正;每次插入时把 updated_at 设为 now(),意味着 RMT 后台去重时最新版本胜出;delete_insert 是主要的正确性路径,RMT 是安全网 |
agg_hourly_zone_trips | ReplacingMergeTree(updated_at) | updated_at | 滚动重算会为相同的 (hour_bucket, zone_id) 组合重新插入聚合;RMT 确保过期聚合在后台合并时被移除 |
dim_taxi_zones | MergeTree() | , | 由 dbt 完全重载意味着每次运行都做原子表交换(完全重建);不可能累积重复;无需去重 |
dim_payment_type | MergeTree() | , | 与 dim_taxi_zones 相同 |
dim_vendor | MergeTree() | , | 与 dim_taxi_zones 相同 |
| 模型 | unique_key | incremental_strategy | 增量过滤条件 | 为什么用这个过滤条件? |
|---|
fact_trips | trip_id | delete_insert | WHERE 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_insert | WHERE pickup_at >= now() - INTERVAL 2 HOUR | 滚动 2 小时窗口强制对边界小时重新聚合,使不完整小时的计数总能被修正;若用 max(pickup_at) 高水位线,边界小时会被永久低估 |
| 模型 | FROM 子句中是否加 FINAL? | 原因 |
|---|
stg_trips | 是,FROM trips_raw FINAL | trips_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 应当与这里的关键决策一致,或者明确记录你为什么做了不同的选择。