03 开通与迁移
用 Terraform 开通 ClickHouse Cloud,按你的方案创建目标表,并用可续跑的 Python 迁移脚本搬迁 5000 万行数据。
起点
模块 02 已完成:migration-plan.md 已填好,完成度检查清单中每个复选框都已勾选,
而且 Snowflake 生产者仍在运行。本模块的 setup.sh
会检查该文件,缺失或不完整时给出警告,但它绝不会阻断,这里没有任何机制阻止你不带它继续,
只有你自己对接下来两个模块的理解会受影响。总共预留大约 60 分钟,其中约 40-50 分钟是
无人值守的数据传输,你可以让它在后台跑着。这里也是 ClickHouse Cloud 试用额度开始消耗的
地方:开通服务并做完本模块大约消耗 $1-2 的试用额度(整个实验总计约 $2-4)。
为什么
本模块是方案变为现实的地方。你在模块 02 写进
migration-plan.md 的每一个决策,每张表用哪个 MergeTree 引擎、从实际查询负载推导出的
ORDER BY 键、Snowflake 专有构造如何翻译,都会在这里直接敲进建表 DDL,而不是重新从头推导。
ClickHouse 没有你可以事后加装的索引:如果一张表里已经躺着 5000 万行数据,才发现 ORDER BY
键选错了,修复手段是完整重载,而不是一句快速的 ALTER。
这也是为什么那道软性关卡即使无法阻止你,仍然很重要。如果你不带完成的方案就跑本
模块,你在机械层面上仍会成功,dbt run 依然会把 fact_trips 建成 ReplacingMergeTree、
迁移脚本依然会搬完 5000 万行,但你不会知道为什么是这个引擎而不是普通的 MergeTree、为什么
sort key 是这个形状,也不知道该如何为模块 04 稍后展示给你的约 6-9 倍基准提速做论证。
下面的决策对齐表把本模块实现的每一个选择映射回它所回答的那道练习表题目,这样你就能在开通
任何资源之前,用它核对自己的方案。
概念:底层原理
目标架构。 Snowflake 通过行程生产者持续写入新行程, 同时一段一次性的 Python 脚本把已有的 5000 万行回填进 ClickHouse,在迁移期间, 这两个系统是并行运行的,而不是一次切换。
在 ClickHouse 一侧,trips_raw 是脚本写入的落地表。随后 dbt 在它之上
构建 staging 视图和分析层的其余部分,本模块创建了这套 schema,但除 trips_raw 之外
尚未填充数据:
图示配色说明:
- 绿色:数据源(切换前与切换后的行程生产者)
- 蓝色:Snowflake 表
- 橙色:dbt 模型与流水线
- 红色:ClickHouse 表与 materialized view
- 青色:Apache Superset 仪表板
- 虚线箭头:切换后的数据流
为什么用 Python 脚本而不是原生连接器。 把数据从 Snowflake 搬到 ClickHouse 有若干种方法。本实验使用 Python 批处理脚本,下面是与各种替代方案对比后的理由:
| 方法 | 工作原理 | 为什么这里不用 |
|---|---|---|
| ClickPipes(Snowflake 数据源) | ClickHouse Cloud 原生连接器,零 ETL、托管界面 | Snowflake 不是受支持的 ClickPipes 数据源。 ClickPipes 支持 Kafka、S3、Kinesis、PostgreSQL CDC、MySQL CDC 和对象存储。 |
| S3 导出 → ClickPipes S3 | COPY INTO @stage 把 Parquet/CSV 导出到 S3;ClickPipes S3 连接器把它加载进 ClickHouse | 需要一个 S3 bucket、一个 IAM 角色、一个 Snowflake stage 和一个 AWS 账号。在任何数据搬迁之前多出约 3 个准备步骤。在生产中可行,但对实验来说基础设施太重。 |
S3 导出 → clickhouse-client | 同样的 S3 导出,但用 INSERT INTO ... SELECT FROM s3(...) 加载 | 同样的 S3 前置条件。还要求合作伙伴手动管理文件分块与可续跑性。 |
| Snowflake → Kafka → ClickHouse | Snowflake CDC 流写入一个 Kafka topic;ClickPipes Kafka 连接器摄取它 | 完整的流式流水线,适合生产中对亚分钟级延迟的要求。Kafka 集群对实验环境而言过于笨重。 |
| Python 脚本(本实验) | snowflake-connector-python 以每批 10 万行的游标方式读取;clickhouse-connect 直接插入 | 除了实验本就需要的包之外,零额外基础设施。通过 --resume 可续跑(以 max(pickup_at) 作为水位线)。有实时进度输出。5000 万行约需 40-50 分钟,速率约 2 万行/秒,对一次性迁移练习来说可以接受。 |
为什么 Python 脚本对本实验是正确选择:
- 不需要 AWS 账号。 基于 S3 的方案需要创建 bucket、配置 IAM 策略, 以及一个 Snowflake 外部 stage,这三个准备步骤和 ClickHouse 毫无关系。
- 自包含。 那两个包(
snowflake-connector-python、clickhouse-connect) 装在与 dbt 相同的 venv 里。没有新服务,没有新凭据。 - 可续跑。
--resume让这个脚本可以安全地中断和重启。ReplacingMergeTree(_synced_at)确保重试时的重复插入会被自动 去重。 - 透明。 合作伙伴可以阅读脚本、理解列映射,并根据自己的 schema 进行调整。 与仅仅完成一遍界面向导相比,这种方式更有助于理解迁移过程。
处理迁移间隔。 迁移脚本运行期间(约 40-50 分钟),Snowflake 生产者一直在 运行。在那段时间窗口内写入 Snowflake 的任何行程都不在 ClickHouse 中。本实验在切换时用两遍法来弥合这个间隔,模块 05 会直接带你走一遍:
- 停止 Snowflake 生产者以冻结数据集。
- 运行
python scripts/02_migrate_trips.py --resume,只有增量行会被 传输(是秒级,不是分钟级)。 - 启动 ClickHouse 生产者。
处理迁移重试的那套 ReplacingMergeTree(_synced_at) 去重机制同样也处理这件事:如果本模块
这次运行与稍后的 --resume 那一遍之间有任何数据行重叠,_synced_at 较晚的那一行胜出。
在生产中你什么时候会选择 S3。 如果数据集超过 5 亿行,或者对整表扫描而言 Snowflake warehouse 的查询成本很显著,那么 S3 导出路径更可取: Snowflake 会并行导出压缩后的 Parquet(比单个游标快得多), 而 ClickHouse 也能并行地从 S3 加载。这里的 Python 脚本方案在实验规模下工作良好。
决策对齐。 下面这张表就是 migration-plan.md 中练习表 1
(引擎选择)、2(sort key)和 3(schema 转换)里的同一份决策清单,
并与本实验实际构建出来的东西做了交叉核对,在开通任何资源之前,先拿它和你自己的方案对比:
| 决策 | 本实验的实现 | 原因 |
|---|---|---|
trips_raw 引擎 | ReplacingMergeTree(_synced_at) | Python 迁移脚本使用批量 INSERT,中断后可能被重试。每次 INSERT 都会设置 _synced_at DateTime DEFAULT now(),因此重试到达的行时间更晚、_synced_at 值更大,在 RMT 去重时较晚的那一行胜出,从而让重试具备幂等性。出于同样的原因,切换后生产者的重试也是安全的。stg_trips 使用 FINAL 查询,以保证每个行程只有一行。 |
fact_trips 引擎 | ReplacingMergeTree(updated_at) | 行程可能被修正(车费调整);updated_at 是版本列 |
agg_hourly_zone_trips 引擎 | ReplacingMergeTree(updated_at) | 滚动重算 = upsert 模式 |
dim_* 表的引擎 | MergeTree() | 每次 dbt 运行都整体重载;没有 upsert |
fact_trips 的 ORDER BY | (toStartOfMonth(pickup_at), pickup_at, trip_id) | Q1-Q7 全都在 pickup_at 上过滤;trip_id 在最细粒度上保证唯一性 |
agg_hourly_zone_trips 的 ORDER BY | (hour_bucket, zone_id) | 两列都出现在所有聚合查询中 |
| VARIANT → | String + JSONExtract* | 保留原始 JSON;提取在查询时发生 |
| QUALIFY → | 用子查询包裹 ROW_NUMBER() | ClickHouse 自 v24.5 起就有原生的 QUALIFY 子句,但这里教子查询形式,因为它可移植到早于该版本或不支持 QUALIFY 的 ClickHouse 版本与 SQL 引擎上 |
| MERGE INTO → | dbt 中的 delete_insert 增量模式 | dbt-clickhouse 惯用的 upsert 策略;避免整表重写 |
步骤 1:开通 ClickHouse 集群
setup.sh 只做一件事:运行 terraform apply 并把连接信息写入
.clickhouse_state。在创建任何资源之前,它还会再次检查 migration-plan.md。
如上文"为什么"一节所述,这项检查只会给出警告,不会中断操作。
cd "$(git rev-parse --show-toplevel)/workshop_public/snowflake_migration_lab/03-migrate-to-clickhouse"
# Configure credentials
cp .env.example .env
vim .env
# Fill in: CLICKHOUSE_ORG_ID, CLICKHOUSE_TOKEN_KEY, CLICKHOUSE_TOKEN_SECRET, CLICKHOUSE_PASSWORD
# Provision
source .env && ./setup.sh.env 在 gitignore 中,绝不要提交它。
预期输出: Terraform 在 大约 2-3 分钟内创建 2 个资源(服务 + IP 访问列表):
Apply complete! Resources: 2 added, 0 changed, 0 destroyed.
Outputs:
clickhouse_host = "abc123xyz.us-east-1.aws.clickhouse.cloud"
clickhouse_port = 8443主机名和端口会被保存到 .clickhouse_state。在任意终端里 source 它即可获得
连接信息:
source .clickhouse_state验证:
curl "https://${CLICKHOUSE_HOST}:${CLICKHOUSE_PORT}/?query=SELECT+1" \
--user "default:${CLICKHOUSE_PASSWORD}"
# Expected: 1步骤 2:创建空表
首先,手动用正确的引擎创建 trips_raw。步骤 3 中迁移脚本会把
数据加载进这张表,它必须已经以 ReplacingMergeTree 存在,这样在任何数据行到达之前
版本列就已经就位。
-- Run in the ClickHouse SQL console (cloud.clickhouse.com -> SQL console)
CREATE TABLE IF NOT EXISTS default.trips_raw (
trip_id String,
vendor_id UInt8,
pickup_at DateTime64(3, 'UTC'),
dropoff_at DateTime64(3, 'UTC'),
passenger_count UInt8,
trip_distance_miles Float32,
pickup_location_id UInt16,
dropoff_location_id UInt16,
payment_type_id UInt8,
rate_code_id UInt8,
store_fwd_flag String,
fare_amount_usd Float32,
extra_amount_usd Float32,
mta_tax_usd Float32,
tip_amount_usd Float32,
tolls_amount_usd Float32,
total_amount_usd Float32,
ingested_at DateTime64(3, 'UTC'),
trip_metadata String,
_synced_at DateTime DEFAULT now()
)
ENGINE = ReplacingMergeTree(_synced_at)
ORDER BY (pickup_at, trip_id);_synced_at 在每次 INSERT 时自动设置。如果迁移脚本被中断
并用 --resume 重跑,同一个 trip_id 可能会短暂存在重复行,RMT
会保留较晚的那一行(_synced_at 更大)。stg_trips 查询 trips_raw FINAL,以在任何
下游模型看到数据之前强制去重。
接下来,灌入区域参考数据。这是静态数据(265 个 NYC TLC 区域),dbt 的
stg_taxi_zones 会把它作为 source 读取。
cd "$(git rev-parse --show-toplevel)/workshop_public/snowflake_migration_lab/03-migrate-to-clickhouse"
source .env && source .clickhouse_state
clickhouse-client --host "${CLICKHOUSE_HOST}" --port 9440 --secure \
--user default --password "${CLICKHOUSE_PASSWORD}" \
--multiquery < scripts/00_seed_zones.sql(或者把 scripts/00_seed_zones.sql 的内容直接粘贴进 ClickHouse SQL
控制台。)
配置 dbt profile。 本项目的 dbt_project.yml 声明了
profile: 'nyc_taxi_ch'。如果 ~/.dbt/profiles.yml 中没有对应的 profile,dbt run
会立刻以 Could not find profile named 'nyc_taxi_ch' 失败。模板位于
workshop_public/snowflake_migration_lab/03-migrate-to-clickhouse/dbt/nyc_taxi_dbt_ch/profiles.yml.example。
模块 01 已经写好了带有 Snowflake 用 nyc_taxi: profile 的 ~/.dbt/profiles.yml,
而且只要 Snowflake 生产者还在运行,模块 01 步骤 4 的刷新循环就会一直用那个 profile
执行查询。不要用 ClickHouse 模板替换那个文件,
用 profiles.yml.example 覆盖它会删掉 nyc_taxi: profile,从而破坏
模块 01 的刷新循环。相反,请打开模板,把它的 nyc_taxi_ch: 块合并
进你现有的 ~/.dbt/profiles.yml,作为与 nyc_taxi: 并列的第二个顶层
profile:
nyc_taxi: # from module 01 — leave this one alone
target: dev
outputs:
dev:
type: snowflake
# ...
nyc_taxi_ch: # add this block
target: dev
outputs:
dev:
type: clickhouse
schema: nyc_taxi_ch
host: "{{ env_var('CLICKHOUSE_HOST') }}"
port: 8443
user: "{{ env_var('CLICKHOUSE_USER', 'default') }}"
password: "{{ env_var('CLICKHOUSE_PASSWORD') }}"
secure: truenyc_taxi_ch: 通过 env_var() 从环境中读取 CLICKHOUSE_HOST、CLICKHOUSE_USER 和
CLICKHOUSE_PASSWORD,因此本模块中任何 dbt 命令之前都必须先 source .env 和
.clickhouse_state,下面的 dbt run 已经这么做了。和
Snowflake profile 一样,~/.dbt/profiles.yml 保存凭据且在 gitignore 中;绝不要提交
它,而这次合并是在一个本就装着一套凭据的文件里再加入第二套。
验证:
cd "$(git rev-parse --show-toplevel)/workshop_public/snowflake_migration_lab/03-migrate-to-clickhouse/dbt/nyc_taxi_dbt_ch"
source "$(git rev-parse --show-toplevel)/workshop_public/snowflake_migration_lab/03-migrate-to-clickhouse/.clickhouse_state"
source "$(git rev-parse --show-toplevel)/workshop_public/snowflake_migration_lab/03-migrate-to-clickhouse/.env"
dbt debug
# Expected: "All checks passed!" — confirms dbt found the nyc_taxi_ch profile and
# connected to ClickHouse然后运行 dbt run 创建分析表和 staging 视图:
cd "$(git rev-parse --show-toplevel)/workshop_public/snowflake_migration_lab/03-migrate-to-clickhouse/dbt/nyc_taxi_dbt_ch"
source "$(git rev-parse --show-toplevel)/workshop_public/snowflake_migration_lab/03-migrate-to-clickhouse/.clickhouse_state"
source "$(git rev-parse --show-toplevel)/workshop_public/snowflake_migration_lab/03-migrate-to-clickhouse/.env"
dbt deps # install packages (first run only)
dbt run # creates analytics tables and staging views; all empty at this point预期: 约 8 个模型在 2 分钟内创建完成(所有表都是空的)。
验证:
# Check analytics tables were created
curl "https://${CLICKHOUSE_HOST}:${CLICKHOUSE_PORT}/?query=SHOW+TABLES+IN+analytics" \
--user "default:${CLICKHOUSE_PASSWORD}"
# Expected: agg_hourly_zone_trips, dim_date, dim_payment_type, dim_vendor, dim_taxi_zones, fact_trips
# Check staging views were created
curl "https://${CLICKHOUSE_HOST}:${CLICKHOUSE_PORT}/?query=SHOW+TABLES+IN+staging" \
--user "default:${CLICKHOUSE_PASSWORD}"
# Expected: stg_trips, stg_taxi_zones
# Check trips_raw exists with the correct engine
curl "https://${CLICKHOUSE_HOST}:${CLICKHOUSE_PORT}/?query=SELECT+engine+FROM+system.tables+WHERE+database%3D%27default%27+AND+name%3D%27trips_raw%27" \
--user "default:${CLICKHOUSE_PASSWORD}"
# Expected: ReplacingMergeTree现在全部六张分析表和两个 staging 视图都已存在,但它们每一个
都还是空的,dbt run 只创建了它们的 schema。这一步之后唯一有数据的表
是 trips_raw,而它此刻也还没有数据;那是下一步的事。
步骤 3:迁移数据
用 Python 批处理迁移脚本,把所有数据行从 Snowflake NYC_TAXI_DB.RAW.TRIPS_RAW 加载进
ClickHouse 的 default.trips_raw。
cd "$(git rev-parse --show-toplevel)/workshop_public/snowflake_migration_lab/03-migrate-to-clickhouse"
source .env && source .clickhouse_state
source .venv/bin/activate
python scripts/02_migrate_trips.py预期输出(5000 万行约需 40-50 分钟):
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━
NYC Taxi Migration: Snowflake -> ClickHouse
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━
Rows to migrate: 50,000,000
Batch size: 100,000
Rows inserted Elapsed ETA Rate
-------------------- ------------ ---------------------- ---------------
100,000 0m 07s 56m 14s remaining 13,945 rows/s
200,000 0m 14s 55m 28s remaining 14,021 rows/s
...这个脚本运行的全程,Snowflake 生产者都在持续写入 TRIPS_RAW,因此
ClickHouse 大致会落后一个本次传输的时长,这个间隔是预期之内的,
将在模块 05 处理,而不是这里。
如果脚本被中断,用 --resume 重跑即可从上一个检查点
继续:
python scripts/02_migrate_trips.py --resume--resume 会从 ClickHouse 读取 max(pickup_at) 并跳过已加载的行,所以
中断并重启这个脚本永远是安全的,你绝不会落到一个不完整、不可恢复的加载状态。
如何确认已完成
cd "$(git rev-parse --show-toplevel)/workshop_public/snowflake_migration_lab/03-migrate-to-clickhouse"
source .env && source .clickhouse_state
# Row count in trips_raw
curl "https://${CLICKHOUSE_HOST}:${CLICKHOUSE_PORT}/?query=SELECT+count()+FROM+default.trips_raw" \
--user "default:${CLICKHOUSE_PASSWORD}"
# Expected: approximately 50000000
# .clickhouse_state was written by setup.sh
ls -la .clickhouse_state
# Service is reachable
curl "https://${CLICKHOUSE_HOST}:${CLICKHOUSE_PORT}/?query=SELECT+1" \
--user "default:${CLICKHOUSE_PASSWORD}"
# Expected: 1结束状态
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)均已就绪。此外,
analytics schema 中还有可刷新的 materialized view
analytics.mv_live_trip_feed,因此共有七个对象。这些对象都由本模块的 dbt run
构建。除了 mv_live_trip_feed,每张
analytics 表都还是空的;而它已经装着 dbt run 构建该视图时产出的那一行快照数据。除此之外,只有 default.trips_raw 有数据:
大约 5000 万行,由步骤 3 中的 Python 迁移脚本搬入。
Snowflake 生产者仍在运行。 本模块中从未停过它,这里也不会停。迁移脚本最后一批之后
写入 Snowflake TRIPS_RAW 的每一条行程,都是 ClickHouse 所没有的数据行,因此
ClickHouse 现在落后 Snowflake 大约一个迁移窗口的时长(约 40-50 分钟,再加上本模块
准备工作所耗的时间)。这个间隔是真实存在的,并且只要生产者还在运行就会不断
扩大。不要在本模块里去弥合它。 模块 05 的切换会刻意弥合它,
用一个受控的两遍步骤,在消除间隔之前先测量它的大小,
现在停掉生产者或重跑迁移脚本,会把模块 05 专门用来演示的那个东西给抹掉。