03 プロビジョニングと移行
Terraform で ClickHouse Cloud をプロビジョニングし、計画からターゲットテーブルを作成し、再開可能な Python マイグレーションスクリプトで5,000万行を移します。
開始チェックポイント
モジュール02が完了していること。migration-plan.md が記入済みで、Completion Checklist の
チェックボックスがすべてチェックされており、Snowflake のプロデューサーがまだ動いていること。この
モジュールの setup.sh はそのファイルをチェックし、欠けている場合や未完成の場合は警告しますが、
決してブロックはしません。それなしで進めることを止めるものはここには何もなく、止めるのは次の
2モジュールに対するあなた自身の理解度だけです。合計で約60分を見込んでください。そのうち約40〜50分は
バックグラウンドに任せて放置できる無人のデータ転送です。ここが ClickHouse Cloud のトライアル支出が
始まる場所でもあります。サービスのプロビジョニングとこのモジュールの実施で、トライアルクレジットの
およそ$1〜2を消費します(ラボ全体では合計で約$2〜4です)。
なぜ必要か
このモジュールは、計画が現実になる場所です。モジュール02で migration-plan.md に書き込んだ
すべての判断 — テーブルごとの MergeTree エンジン、実際のクエリワークロードから導いた ORDER BY
キー、Snowflake 固有の構文をどう変換するか — が、ここでゼロから再導出されるのではなく、そのまま
テーブルの DDL に打ち込まれます。ClickHouse には、あとから後付けできるインデックスがありません。
5,000万行がテーブルに収まったあとで ORDER BY キーが間違っていたと判明した場合、対処は手軽な
ALTER ではなく、フルリロードです。
だからこそ、ソフトなゲートはあなたを止められないにもかかわらず重要です。完成した計画なしにこの
モジュールを実行しても、機械的には成功します。dbt run は fact_trips を ReplacingMergeTree
として作りますし、マイグレーションスクリプトは5,000万行を移します。しかし、なぜ素の MergeTree
ではなくそのエンジンなのか、なぜ sort key がその形なのか、モジュール04であとから見ることになる
約6〜9倍のベンチマーク高速化をどう正当化するのかは分かりません。下の判断の対応表は、この
モジュールが実装するすべての選択を、それが答えているワークシートの問いに紐付けています。何かを
プロビジョニングする前に、自分の計画をこれと照合してください。
概念 — 内部の仕組み
ターゲットのアーキテクチャ。 Snowflake は trip プロデューサー経由で新しい trip を書き込み 続け、その間に一度きりの Python スクリプトが既存の5,000万行を ClickHouse にバックフィルします。 2つのシステムはマイグレーションの期間中は並走し、カットオーバーではありません。
ClickHouse 側では、trips_raw がスクリプトの書き込み先となるランディングテーブルです。その後
dbt が、その上に staging view と analytics レイヤーの残りをビルドします。このモジュールが作るのは
このスキーマですが、trips_raw を超える範囲にはまだデータを投入しません。
図の色の凡例:
- 緑 — データソース(trip プロデューサー、カットオーバー前と後)
- 青 — Snowflake のテーブル
- オレンジ — dbt のモデルとパイプライン
- 赤 — ClickHouse のテーブルと materialized view
- シアン — Apache Superset のダッシュボード
- 破線の矢印 — カットオーバー後のフロー
ネイティブコネクターではなく Python スクリプトを使う理由。 Snowflake から ClickHouse に データを移す方法はいくつも存在します。このラボは Python のバッチスクリプトを使います。代替手段と 比べた理由は次のとおりです。
| 方法 | 仕組み | ここで使わない理由 |
|---|---|---|
| ClickPipes(Snowflake ソース) | ClickHouse Cloud のネイティブコネクター — ゼロ ETL、マネージドな UI | Snowflake は ClickPipes のサポート対象ソースではありません。 ClickPipes がサポートするのは Kafka、S3、Kinesis、PostgreSQL CDC、MySQL CDC、オブジェクトストレージです。 |
| S3 エクスポート → ClickPipes S3 | COPY INTO @stage で Parquet/CSV を S3 にエクスポートし、ClickPipes の S3 コネクターが ClickHouse に読み込む | S3 バケット、IAM ロール、Snowflake のステージ、AWS アカウントが必要です。データが1行も動く前にセットアップ手順が約3つ増えます。本番では有効ですが、ラボにはインフラが多すぎます。 |
S3 エクスポート → clickhouse-client | 同じ S3 エクスポートを、INSERT INTO ... SELECT FROM s3(...) で読み込む | S3 の前提条件は同じです。加えて、ファイルのチャンク分割と再開可能性をパートナーが手動で管理する必要があります。 |
| Snowflake → Kafka → ClickHouse | Snowflake の CDC ストリームが Kafka のトピックに供給し、ClickPipes の Kafka コネクターが取り込む | フルのストリーミングパイプラインで、本番で分未満のレイテンシー要件があるなら適切です。Kafka クラスターはラボ環境には重すぎます。 |
| Python スクリプト(このラボ) | snowflake-connector-python が10万行のカーソルバッチで読み取り、clickhouse-connect が直接挿入する | ラボがすでに必要とするパッケージ以外、追加インフラはゼロです。--resume(max(pickup_at) のウォーターマーク)で再開可能です。進捗をリアルタイムに出力します。約20K rows/s で5,000万行に約40〜50分 — 一度きりのマイグレーション演習としては許容できます。 |
このラボで Python スクリプトが正しい選択である理由:
- AWS アカウントが不要。 S3 ベースのアプローチには、バケット作成、IAM ポリシー、Snowflake の 外部ステージが必要です。ClickHouse とは何の関係もない3つのセットアップ手順です。
- 自己完結している。 2つのパッケージ(
snowflake-connector-python、clickhouse-connect)は dbt と同じ venv にインストールされます。新しいサービスも、新しい認証情報も要りません。 - 再開可能。
--resumeにより、スクリプトは中断して再開しても安全です。ReplacingMergeTree(_synced_at)が、リトライ時の重複挿入を自動的に重複排除します。 - 透明性がある。 パートナーはスクリプトを読み、カラムのマッピングを理解し、自分のスキーマ向けに 適応させられます。UI のウィザードをクリックしていくよりも教育的です。
マイグレーションのギャップの扱い。 マイグレーションスクリプトの実行中(約40〜50分)も Snowflake のプロデューサーは動き続けます。その間に Snowflake に書き込まれた trip は ClickHouse には ありません。このラボはそのギャップを、カットオーバー時の2パス方式で閉じます。モジュール05で実際に 手を動かして進めます。
- Snowflake のプロデューサーを停止してデータセットを凍結する。
python scripts/02_migrate_trips.py --resumeを実行する — 差分の行だけが転送されます (分ではなく秒単位です)。- ClickHouse のプロデューサーを起動する。
マイグレーションのリトライを扱うのと同じ ReplacingMergeTree(_synced_at) の重複排除が、これも
扱います。このモジュールの実行と、あとの --resume パスの間で行が重複しても、_synced_at が
新しい方が勝ちます。
本番で S3 を選ぶ場合。 データセットが5億行を超える場合、あるいは Snowflake ウェアハウスでの フルテーブルスキャンのクエリコストが無視できない場合は、S3 エクスポートの経路が望ましいです。 Snowflake は圧縮された Parquet を並列でエクスポートでき(単一カーソルよりずっと高速です)、 ClickHouse も S3 から並列で読み込めます。ここでの Python スクリプト方式は、ラボの規模には よく合っています。
判断の対応。 下の表は、migration-plan.md のワークシート1(エンジン選定)、2(sort key)、
3(スキーマ変換)と同じ判断リストを、このラボが実際に構築するものと突き合わせたものです。何かを
プロビジョニングする前に、自分の計画と比べてください。
| 判断 | このラボの実装 | 理由 |
|---|---|---|
trips_raw のエンジン | ReplacingMergeTree(_synced_at) | Python のマイグレーションスクリプトはバッチ INSERT を使い、中断した場合にリトライされることがあります。_synced_at DateTime DEFAULT now() はすべての INSERT で設定されるため、リトライされた行はあとから到着し、より大きな _synced_at の値を持ちます。RMT の重複排除ではあとの行が勝つので、リトライは冪等になります。同じ理由から、カットオーバー後のプロデューサーのリトライも安全です。stg_trips は FINAL を付けてクエリし、trip ごとに1行であることを保証します。 |
fact_trips のエンジン | ReplacingMergeTree(updated_at) | trip は訂正されることがあります(運賃の調整)。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 がやることは1つです。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 に対してクエリし、
下流のモデルがデータを見る前に重複排除を強制します。
次に、zone のリファレンスデータをシード投入します。これは静的なデータ(NYC TLC の265の zone)で、
dbt の stg_taxi_zones がソースとして読み取ります。
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 のプロファイルを設定する。 このプロジェクトの dbt_project.yml は
profile: 'nyc_taxi_ch' を宣言しています。~/.dbt/profiles.yml に対応するプロファイルがないと、
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はすでに ~/.dbt/profiles.yml に Snowflake 用の nyc_taxi: プロファイルを書き込んで
おり、その手順4のリフレッシュループは、Snowflake のプロデューサーが動いている間ずっとその
プロファイルに対してクエリを続けます。そのファイルを ClickHouse のテンプレートで置き換えては
いけません。 profiles.yml.example で上書きすると nyc_taxi: プロファイルが消え、モジュール01の
リフレッシュループが壊れます。代わりにテンプレートを開き、その nyc_taxi_ch: ブロックを、
nyc_taxi: と並ぶ2つ目のトップレベルプロファイルとして既存の ~/.dbt/profiles.yml にマージして
ください。
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 コマンドを実行する前に .env と
.clickhouse_state を source しておく必要があります。下の dbt run はすでにそうしています。
Snowflake のプロファイルと同様、~/.dbt/profiles.yml は認証情報を保持し、gitignore されています。
決してコミットしないでください。またこのマージは、すでに1組の認証情報を持つファイルに2組目を
追加するものです。
検証:
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 を実行し、analytics のテーブルと staging の view を作成します。
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期待される結果: 2分以内に約8つのモデルが作成されます(テーブルはすべて空)。
検証:
# 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: ReplacingMergeTree6つの analytics テーブルと2つの staging view はこれで存在しますが、そのすべてがまだ空です。
dbt run はスキーマを作っただけです。この手順のあとにデータがあるテーブルは 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期待される出力(5,000万行で約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 のサービスが稼働し到達可能で、.clickhouse_state が setup.sh によって
CLICKHOUSE_HOST と CLICKHOUSE_PORT とともにディスクに書き出されています。ターゲットの
テーブルと staging view はすべて存在します。default.trips_raw、2つの staging view(stg_trips、
stg_taxi_zones)、6つの analytics テーブル(fact_trips、agg_hourly_zone_trips、
dim_taxi_zones、dim_payment_type、dim_vendor、dim_date)、そしてリフレッシュ可能な
materialized view analytics.mv_live_trip_feed — 合計で analytics スキーマに7つのオブジェクトが、
このモジュールの dbt run によって作られました。analytics のテーブルはすべてまだ空ですが、
mv_live_trip_feed は例外で、dbt run が view をビルドしたときに生成したスナップショット1行を
すでに保持しています。それ以外でデータがあるのは default.trips_raw のみで、手順3の Python
マイグレーションスクリプトが移したおよそ5,000万行が入っています。
Snowflake のプロデューサーはまだ動いています。 このモジュールで止めたことはなく、ここでも
止めません。マイグレーションスクリプトの最後のバッチ以降に Snowflake の TRIPS_RAW に書き込まれた
trip はすべて、ClickHouse が持っていない行です。したがって ClickHouse はいま、おおよそ
マイグレーションのウィンドウの長さ(約40〜50分に、このモジュールのセットアップにかかった時間を
加えたもの)だけ Snowflake より遅れています。そのギャップは実在し、プロデューサーが動いている
限り広がり続けます。このモジュールでは閉じないでください。 モジュール05のカットオーバーが、
ギャップを解消する前にその大きさを測定する制御された2パスの手順で、意図的に閉じます。いま
プロデューサーを止めたりマイグレーションスクリプトを再実行したりすると、モジュール05が示すために
作られているまさにそのものが消えてしまいます。