01 源环境
开通一个映射真实客户部署的 Snowflake 环境,5000 万行数据、一条 dbt Medallion 流水线、一个实时行程生产者以及三个 Superset 仪表板。
起点
模块 00 已完成:工具链已安装、两个云试用账号均已可用、仓库已克隆、dbt-snowflake 虚拟环境
已构建。本模块约需 45 分钟,大致消耗 2-4 个 Snowflake 积分。
为什么
你无法基于一个玩具级源端去规划迁移。一张只有寥寥几行的扁平表会让你绕过所有让真实迁移变难
的决策。本模块转而搭建一个真实客户部署的形态:一个存放半结构化 JSON 的 VARIANT 列、一条
CDC 流、定时任务、一条增量 MERGE 流水线,以及在这一切之上做读取的 BI 层。其中每一项都会
在模块 02 变成一个具体的迁移决策,本模块存在的意义,就是让那个决策出现时,你有真实的东西
可以指着看,而不是一个抽象概念。
概念:底层原理
基础设施(Terraform)。 运行 setup.sh 会开通:
- Warehouse:
TRANSFORM_WH(SMALL,用于 ELT)和ANALYTICS_WH(MEDIUM,用于 BI), 外加一个上限为每月 50 积分的资源监控器(ANALYTICS_WH_MONITOR)。 - 数据库:
NYC_TAXI_DB,包含三个 schema:RAW、STAGING、ANALYTICS。 - 角色:
TRANSFORMER_ROLE、ANALYST_ROLE、DBT_ROLE、LOADER_ROLE。
Medallion 形态。 数据在 NYC_TAXI_DB 内部流经三层:
- RAW:
TRIPS_RAW(5000 万行合成行程数据,其中包含一个模拟应用遥测的TRIP_METADATAVARIANT列,这就是 JSON 迁移挑战)以及维度表 (DIM_TAXI_ZONES、DIM_PAYMENT_TYPE、DIM_VENDOR)。 - STAGING:用于清洗类型并展平
VARIANT列的 dbt 视图。 - ANALYTICS:dbt 表和增量模型:
fact_trips(5000 万行,MERGE策略)、四张维度表,以及agg_hourly_zone_trips(一个增量 聚合)。
有两个 Snowflake 对象独立于 dbt、自行驱动这条流水线:
TRIPS_CDC_STREAM:TRIPS_RAW上的一条变更数据捕获流。CDC_CONSUME_TASK:每 5 分钟读取该流(运行在RAW中,在 setup 期间被 恢复),以及HOURLY_AGG_TASK,它每小时刷新一次小时级聚合 (运行在STAGING中,在 dbt 构建之后被恢复)。
Superset。 三个仪表板都通过 ANALYTICS_WH 从 ANALYTICS schema 读取,它们都不直接
访问 RAW 或 STAGING。这条读取路径就是你稍后要在 ClickHouse 一侧
复现的东西。
步骤 1:配置凭据
cd "$(git rev-parse --show-toplevel)/workshop_public/snowflake_migration_lab/01-setup-snowflake"
cp .env.example .env
# Edit .env with your Snowflake credentials
cp dbt/nyc_taxi_dbt/profiles.yml.example ~/.dbt/profiles.yml
# Edit ~/.dbt/profiles.yml with your account details.env 和 ~/.dbt/profiles.yml 都在 gitignore 中,它们保存你的 Snowflake 账号、
用户名和密码。绝不要提交这两个文件中的任何一个。
步骤 2:运行 setup
cd "$(git rev-parse --show-toplevel)/workshop_public/snowflake_migration_lab/01-setup-snowflake"
source .env && ./setup.sh预计这会花费 5-10 分钟,其中大部分时间用于通过 TABLE(GENERATOR) 生成 5000 万行
合成行程数据。setup.sh 一次性完成:开通 Terraform 基础设施、灌入 TRIPS_RAW 数据、
运行 dbt 构建,并启动 Docker Compose(行程生产者和 Superset)。
步骤 3:启动生产者与 Superset
setup.sh 会用正确的环境变量启动 Docker Compose、在 Superset 中注册
Snowflake 连接,并自动导入全部三个仪表板。
superset/dashboards/ 中已提交的仪表板 ZIP 文件,其 sqlalchemy_uri 已被
脱敏为占位符(LAB_USER、MYORG-MYACCOUNT)。自动导入会用你的 .env 重新写入该
URI,因此在 setup.sh 运行时这一步是透明的。如果你改为通过 Superset 界面手动导入
ZIP,它创建的连接会使用那些占位符从而连不上,之后请编辑该连接,指向你真实的
Snowflake 账号。
如果你需要手动重启 Superset:
cd "$(git rev-parse --show-toplevel)/workshop_public/snowflake_migration_lab/01-setup-snowflake/superset"
docker-compose --env-file ../.env up -d--env-file ../.env 参数会从上级目录加载环境变量。
三个仪表板分别是 Operations Command Center、Executive Weekly Report 和 Driver & Quality Analytics(后者故意做得很慢,它是后续 ClickHouse 基准测试的 对象)。完整的仪表板搭建过程,数据源、图表、筛选器,请见 Snowflake 上的 Superset。
步骤 4:让 dbt 保持最新
行程生产者持续以约每分钟 60 条行程的速度向 TRIPS_RAW 插入数据。要在你工作期间让
fact_trips 和 agg_hourly_zone_trips 保持最新,请在另一个终端里运行 dbt 刷新
循环:
cd "$(git rev-parse --show-toplevel)/workshop_public/snowflake_migration_lab/01-setup-snowflake"
# Default: refresh every 5 minutes (auto-sources .env)
./scripts/run_dbt.sh
# Custom interval
./scripts/run_dbt.sh --interval 15m
# Run once and exit
./scripts/run_dbt.sh --once
# Include dbt tests after each run
./scripts/run_dbt.sh --test| 参数 | 效果 |
|---|---|
--interval <n> | 两次运行之间的间隔:30s、5m、1h,或纯秒数(默认:5m) |
--once | 只刷新一次然后退出 |
--test | 在每次 dbt run 之后运行 dbt test |
该脚本始终以增量方式运行,它从不执行 --full-refresh,因此生产者插入的数据行会被
保留。随时按 Ctrl-C 停止它;在本实验的剩余部分中,请让它在自己的终端里持续运行,模块 03
仍然依赖它。
步骤 5:浏览查询库
cd "$(git rev-parse --show-toplevel)/workshop_public/snowflake_migration_lab/01-setup-snowflake"查询目录(workshop_public/snowflake_migration_lab/01-setup-snowflake/queries/)
中有七个带注释的 SQL 文件。每一个都能对你刚搭好的 Snowflake 环境运行,并且每一个都刻意
带有一项模块 02 将会翻译到 ClickHouse 的迁移挑战:
| 查询 | 构造 | 迁移挑战 |
|---|---|---|
| Q1 | DATE_TRUNC、DATEADD | 细微的语法差异 |
| Q2 | 窗口函数 ROWS BETWEEN | 在 ClickHouse 中几乎完全一致 |
| Q3 | QUALIFY | ClickHouse 自 v24.5 起原生支持,此处为了可移植性仍改写为子查询 |
| Q4 | LATERAL FLATTEN | 没有等价物,使用 JSONExtract 或预先展平 |
| Q5 | VARIANT 冒号路径 | 替换为 JSONExtractFloat/JSONExtractString |
| Q6 | MERGE INTO | 没有等价物,使用 ReplacingMergeTree |
| Q7 | Snowflake Streams | 切换时退役,实时写入经生产者直接进入 ClickHouse |
在继续之前,打开每个文件并针对你的 Snowflake 环境运行一次。每个查询里的注释块已经勾勒出 ClickHouse 的等价写法,模块 02 才是你真正动手编写并运行那一侧的地方。
如何确认已完成
cd "$(git rev-parse --show-toplevel)/workshop_public/snowflake_migration_lab/01-setup-snowflake"
source .env && ./scripts/verify_environment.sh它会检查:
- 数据库与 schema:
NYC_TAXI_DB存在,且包含RAW、STAGING、ANALYTICS。 - 表与数据:
TRIPS_RAW约有 5000 万行、FACT_TRIPS已填充数据、各维度表 均存在。 - CDC 流:
TRIPS_RAW上存在TRIPS_CDC_STREAM。 - 定时任务:
CDC_CONSUME_TASK和HOURLY_AGG_TASK处于started状态。 - CDC 活动:这些任务近期已执行过。
- 生产者数据流:行程生产者正在持续插入数据。
- Superset:BI 仪表板可通过
http://localhost:8088访问。
如果你更愿意手动检查,用 ACCOUNTADMIN 角色执行
SHOW TASKS LIKE '%TASK' IN DATABASE NYC_TAXI_DB;(任务归该角色所有)即可确认两个任务都在运行。
收尾
在接下来的几个模块里你会反复回到这个环境,而完整跑一次 ./setup.sh 要花掉 5-10 分钟,
你不会想在每次改一个 Terraform 文件或 dbt 模型时都付这个代价。setup.sh 正是为此提供了
参数:
| 参数 | 何时使用 |
|---|---|
| (无) | 首次运行。开通所有资源并生成 5000 万行合成数据(总计约 12 分钟)。 |
--skip-seed | 基础设施已存在且 TRIPS_RAW 已有数据。跳过合成数据生成(节省约 8 分钟)。 |
--skip-dbt | Snowflake 对象已存在,但你不需要重跑 dbt 转换(例如测试 Terraform 变更)。 |
--skip-superset | Docker 没在运行,或者你暂时不需要 BI 层。 |
--full-refresh | 强制 dbt 从零重建所有增量模型(例如 schema 变更之后)。 |
参数可以组合使用。两种常见组合:
# Re-run after a Terraform or SQL change — skip the ~10 min data load
./setup.sh --skip-seed
# Iterate on dbt models only — skip everything else
./setup.sh --skip-seed --skip-superset成本说明。 数据灌入约需 12 分钟、消耗 2 积分(约 $6),一次完整的 dbt 构建约需 8 分钟、消耗 1.5 积分(约 $5)。一次 8 小时的合作伙伴实验课程会再增加 大约 12 积分(约 $36),warehouse 在空闲时会自动挂起,因此在课程之间不再产生费用。 每位合作伙伴每天合计大约 16 积分,约 $47。
结束状态
Snowflake 已经就绪:NYC_TAXI_DB 已完整构建、CDC 流和两个定时任务都在运行、行程生产者
正以约每分钟 60 条行程的速度写入 TRIPS_RAW,三个 Superset 仪表板都可在
http://localhost:8088 访问。
让生产者继续运行。 不要停止 Docker Compose 栈,也不要运行
./teardown.sh,模块 02 到 05 都依赖这个环境保持可用,而模块 05 的切换步骤要测量的正是
迁移期间生产者在 Snowflake 与 ClickHouse 之间制造的确切间隔。现在就拆除资源会让
后续课程会以一种很难追溯到这一步的方式失败。资源拆除在模块 05 的末尾讲解,
不在这里。