Snowflake MigrationClickHouse Workshops

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つのシステムはマイグレーションの期間中は並走し、カットオーバーではありません。

マイグレーションのデータフロー: Snowflake の trip プロデューサーが TRIPS_RAW に書き込み続ける一方、一度きりの Python スクリプトが5,000万行を10万行のバッチで ClickHouse Cloud に移す

ClickHouse 側では、trips_raw がスクリプトの書き込み先となるランディングテーブルです。その後 dbt が、その上に staging view と analytics レイヤーの残りをビルドします。このモジュールが作るのは このスキーマですが、trips_raw を超える範囲にはまだデータを投入しません。

ClickHouse のターゲット側: ReplacingMergeTree 上の trips_raw が dbt でビルドされた staging view、fact テーブルとディメンションテーブル、時間別集計、zone ディクショナリに供給する

図の色の凡例:

  • 緑 — データソース(trip プロデューサー、カットオーバー前と後)
  • 青 — Snowflake のテーブル
  • オレンジ — dbt のモデルとパイプライン
  • 赤 — ClickHouse のテーブルと materialized view
  • シアン — Apache Superset のダッシュボード
  • 破線の矢印 — カットオーバー後のフロー

ネイティブコネクターではなく Python スクリプトを使う理由。 Snowflake から ClickHouse に データを移す方法はいくつも存在します。このラボは Python のバッチスクリプトを使います。代替手段と 比べた理由は次のとおりです。

方法仕組みここで使わない理由
ClickPipes(Snowflake ソース)ClickHouse Cloud のネイティブコネクター — ゼロ ETL、マネージドな UISnowflake は ClickPipes のサポート対象ソースではありません。 ClickPipes がサポートするのは Kafka、S3、Kinesis、PostgreSQL CDC、MySQL CDC、オブジェクトストレージです。
S3 エクスポート → ClickPipes S3COPY 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 → ClickHouseSnowflake の 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で実際に 手を動かして進めます。

  1. Snowflake のプロデューサーを停止してデータセットを凍結する。
  2. python scripts/02_migrate_trips.py --resume を実行する — 差分の行だけが転送されます (分ではなく秒単位です)。
  3. 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: true

nyc_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: ReplacingMergeTree

6つの 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が示すために 作られているまさにそのものが消えてしまいます。

このページの内容

Track your progress?

Optional. We email a link to confirm your address; progress records once you open it.

Please use your work email address, not a personal one.

Progress tracking also requires accepting the current Terms of Service in Privacy settings.

JA