Snowflake MigrationClickHouse Workshops

ワークト例: 完成した計画

NYC タクシーのワークロードについて記入済みの移行計画。自分で書き上げたあとに比較するためのもの。

これは 5 つのワークシートすべて (1、 2、 3、 4、 5) を NYC Taxi のワークロードに適用した完全な解答です。次の用途に使ってください:

  • 各セクションを終えたあと、自分のワークシートの答えを確認する
  • Part 3 が実装している判断の背後にある理由を理解する
  • 自分が別の選択をした場合に、Part 3 の Decision Alignment 表と比較する

これは解答例です — 自分の計画としてここに記入しないでください。代わりに migration-plan.md に記入してください。


完了チェックリスト

  • エンジン選択: 完了
  • sort key 設計: 完了
  • スキーマ変換: 完了
  • 移行ウェーブ計画: 完了
  • dbt モデル設計: 完了

セクション 1: プロファイルのサマリ

メトリクス値
テーブル総数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)
Stream1 (TRIPS_RAW 上の TRIPS_CDC_STREAM)
Task2 (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)

セクション 2: オブジェクト一覧

オブジェクト種別スキーマ行数複雑度グレード備考
trips_rawテーブルraw約 50MB_synced_at をバージョンカラムに持つ RMT。バルクと CDC の重複により重複排除が必要。stg_trips は FINAL を使う必要がある
stg_tripsdbt ビューstaging—BTRIP_METADATA に対する JSONExtract。JSON パスのテストが必要
stg_taxi_zonesdbt ビューstaging—Aパススルー。ごく簡単
int_trips_enricheddbt Ephemeralstaging—ACTE。SQL の差異は親モデル側で処理される
fact_tripsdbt Incrementalanalytics約 50MCRMT エンジン。delete_insert。QUALIFY の書き換え。FINAL が必要
agg_hourly_zone_tripsdbt Incrementalanalytics約 140KBRMT。2 時間のローリング再計算ウィンドウ。パーティション境界を慎重にテストすること
dim_taxi_zonesdbt テーブルanalytics265A静的な参照データ。全件リロード。ごく簡単
dim_payment_typedbt テーブルanalytics6A静的な参照データ。ごく簡単
dim_vendordbt テーブルanalytics3A静的な参照データ。ごく簡単
taxi_zones_dictディクショナリanalytics265BClickHouse 固有の構文。クエリ時に dictGet()
mv_hourly_revenueリフレッシュ可能な MVanalytics—BREFRESH EVERY 構文。原子的な置き換えを検証すること
TRIPS_CDC_STREAM / CDC_CONSUME_TASKSnowflake Stream + Task——DClickHouse に相当機能なし。Part 3 ではプロデューサーの直接カットオーバーで置き換える

セクション 3: エンジン選択の判断

テーブルエンジンバージョンカラム理由
trips_rawReplacingMergeTree(_synced_at)_synced_atPython の移行スクリプト (scripts/02_migrate_trips.py) はバッチをリトライして同じ trip_id を再挿入する可能性があります。カットオーバー後も、ライブのプロデューサーが一時的な障害でリトライすることがあります。_synced_at DateTime DEFAULT now() は INSERT 時に自動設定されるため、後のリトライはより大きなタイムスタンプを持ち、RMT は最新の書き込みを残します。stg_trips は FINAL でクエリし、下流のモデルが動く前に重複排除を強制します。
fact_tripsReplacingMergeTree(updated_at)updated_atトリップは訂正されえます (運賃の調整、ステータス変更)。同じ trip_id が更新後の値で再挿入されます。updated_at は訂正ごとに単調増加するため、RMT の重複排除では大きい値が勝ちます。常に FINAL でクエリしてください。
agg_hourly_zone_tripsReplacingMergeTree(updated_at)updated_atdbt は直近 2 時間を再計算して再挿入します。RMT なしでは古い集計と新しい集計が蓄積し、二重カウントになります。dbt 実行ごとに updated_at を now() に設定することで、最新の値が勝つようにします。
dim_taxi_zonesMergeTree()—dbt による全件リロード (原子的なテーブル入れ替え (全件再構築))。重複が蓄積することはありません。重複排除は不要です。
dim_payment_typeMergeTree()—同様 — 全件リロード。
dim_vendorMergeTree()—同様 — 全件リロード。
mv_hourly_revenueMergeTree()—REFRESHABLE MV は REFRESH ごとに結果セット全体を原子的に置き換えます。upsert はありません。

セクション 4: sort key 設計

テーブルORDER BY理由
trips_raw(pickup_at, trip_id)時間範囲スキャンはまず pickup_at で絞り込みます。trip_id は RMT の重複排除キーであり、RMT がどの行が重複かを判断できるよう ORDER BY に含める必要があります。分析的な範囲スキャンが支配的なので pickup_at を先頭に、trip_id はカーディナリティが高く一意性の判別子としてしか働かないので最後に置きます。
fact_trips(toStartOfMonth(pickup_at), pickup_at, trip_id)7 つの分析クエリすべてが pickup_at で絞り込みます。月のプレフィックスは暦月のデータを隣接ブロックにまとめ、PARTITION BY を追加せずに月次集計での粗いブロックスキップを可能にします。ブロックスキップを妨げずに RMT の一意性を確保するため、trip_id を最後に置きます。
agg_hourly_zone_trips(hour_bucket, zone_id)Q6 (およびすべての集計クエリ) は hour_bucket と zone_id で絞り込みます。hour_bucket の異なる値は約 35K、zone_id は 265 です。時間範囲スキャンが主要なアクセスパターンなので hour_bucket を先頭に、二次的な絞り込み用に zone_id を 2 番目に置きます。
dim_taxi_zones(location_id)265 行 = 1 グラニュール。ORDER BY はパフォーマンスに関係しません。結合キーである location_id は慣習的で、可読性にも役立ちます。

セクション 5: スキーマ変換のメモ

カラムSnowflake の型ClickHouse の型判断の根拠
TRIP_METADATAVARIANTString生の JSON をそのまま保持します。JSONExtract* はクエリ時に任意のパスを扱えます。Map(String,String) ではネスト構造が失われ、Tuple は固定スキーマを要求します。任意の JSON には String が安全な選択です。
PICKUP_DATETIME / PICKUP_ATTIMESTAMP_NTZ(9)DateTime64(3, 'UTC')トリップのタイムスタンプにはミリ秒精度で十分です。ナノ秒 (9) は過剰です。'UTC' はタイムゾーンを明示し、時間範囲の集計での DST に起因する意外な挙動を避けます。
PICKUP_LOCATION_IDINTEGERUInt16値は 1–265。UInt8 の最大値は 255 (小さすぎます)。UInt16 の最大値は 65535 (適切)。Int32 の 4 バイトに対し 2 バイト — 5,000 万行、1 カラムあたり非圧縮で約 95MB 節約します。
VENDOR_IDINTEGERUInt8値は 1–3。UInt8 の最大値は 255 — 適切です。1 行あたり 1 バイト。
DRIVER_RATINGFLOATNullable(Float32)しばしば NULL です (すべてのトリップに評価があるわけではありません)。Nullable は正しい null セマンティクスを保ちます。1.0–5.0 の範囲には Float32 で十分です。Float64 は意味のある精度を加えずにストレージを浪費します。
UPDATED_ATTIMESTAMP_NTZ(9)DateTime64(3, 'UTC')ReplacingMergeTree のバージョンカラム。DateTime ではなく DateTime64 を使う必要があります — 秒精度では同一秒内の 2 回の訂正が非決定的になります。ミリ秒精度により重複排除の順序が正しくなります。

必要な関数の変換

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::FLOATJSONExtractFloat(trip_metadata, 'driver', 'rating')
TRIP_METADATA:surge_multiplier::FLOATJSONExtractFloat(trip_metadata, 'surge_multiplier')
QUALIFY ROW_NUMBER() OVER (PARTITION BY pickup_location_id ORDER BY fare_amount DESC) <= 10SELECT ... FROM (SELECT ..., ROW_NUMBER() OVER (...) AS rn FROM ...) WHERE rn <= 10
MERGE INTO fact_trips ... WHEN MATCHED THEN UPDATEdbt の delete_insert インクリメンタル — キーが一致する行を DELETE し、その後すべての新しい行を INSERT

セクション 6: 移行ウェーブ

ウェーブオブジェクト依存関係備考
ウェーブ 0trips_raw (スキーマ)、dim_taxi_zones、dim_payment_type、dim_vendorなしdbt が空のテーブルを作成します。ディメンションテーブルは静的な参照データから直ちに投入されます (トリップへの依存なし)。実行: dbt run --select trips_raw dim_*
ウェーブ 1Python によるバルクロード (scripts/02_migrate_trips.py)ウェーブ 0 (trips_raw のスキーマが存在していること)Snowflake の TRIPS_RAW から 5,000 万行。--resume で再開可能。scripts/01_verify_migration.sh で行数を検証。約 40〜50 分。
ウェーブ 2stg_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 がその上に構築されます。
ウェーブ 3taxi_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 を投入します。

リスク登録簿 (グレード C/D のオブジェクト)

オブジェクトリスク検証方法
fact_tripsFINAL なしのクエリはマージ遅延の間に過大カウントします。対象外のパーティションを削除しないよう、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_trips2 時間のローリング再計算ウィンドウは削除範囲を正しく限定する必要があります。広すぎると古い集計が削除され、狭すぎると陳腐な集計が残ります。特定の (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) が重複挿入をべき等に処理します。

セクション 7: 既知の方言ギャップ

  • QUALIFY — 影響: Q3 (queries/q03_top_trips_qualify.sql)
  • VARIANT のコロンパス — 影響: Q4、Q5 (TRIP_METADATA の JSON アクセス)
  • LATERAL FLATTEN — このワークロードでは未使用。VARIANT は FLATTEN ではなくコロンパスでアクセスしています
  • MERGE INTO — 影響: dbt のインクリメンタルモデル (fact_trips, agg_hourly_zone_trips)
  • Snowflake Streams → プロデューサーのカットオーバー (カットオーバー後、ライブの書き込みは直接 ClickHouse へ)
  • 日付関数の差異 — 影響: Q1 (DATE_TRUNC)、Q3 (DATEADD)、Q4 (DATEDIFF)

セクション 8: 移行戦略

データ移動: Python の移行スクリプト (scripts/02_migrate_trips.py)

オブジェクトストレージ経由や ClickPipes ではなく Python スクリプトを使う理由:

  • remoteSecure() は 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

append や merge 戦略ではなく delete_insert を使う理由:

  • append は既存の行に触れずに新しい行を挿入します。行が更新されうる fact_trips では重複を生みます。不適切です。
  • merge (利用できる場合) は Snowflake の MERGE INTO に最も近いですが、dbt-clickhouse の merge 戦略は ReplacingMergeTree との組み合わせに制約があり、推奨されるアプローチではありません。
  • delete_insert は到着したバッチのキー範囲の行を削除し、その後すべての新しい行を挿入します。これはべき等で (再実行しても同じ結果になり)、挿入と更新の両方を扱え、ReplacingMergeTree と正しく動作します。upsert パターンに対する dbt-clickhouse コミュニティの標準的な推奨です。

セクション 9: カットオーバーの判定基準

基準しきい値測定方法
行数の一致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 の出力を手動で比較


セクション 10: dbt モデル設計

マテリアライゼーションの選択

モデルマテリアライゼーション理由
stg_tripsviewtrips_raw を読み取ってクリーンにする。このモデル自体に更新はない。ストレージコストがゼロ。常に現在のソースの状態を反映する
stg_taxi_zonesview同様 — ソーステーブルのパススルーなクリーンアップ
int_trips_enrichedephemeralfact_trips からのみ使われる純粋な結合ロジック。CTE としてインライン展開することで冗長な物理テーブルを避ける。直接クエリするモデルはない
fact_tripsincrementalトリップは事後に訂正されうる。実行ごとに新規行と更新行のみを処理すべき
agg_hourly_zone_tripsincremental2 時間のローリング再計算はインクリメンタルなパターン — 5,000 万行すべてではなく最近の行を処理する
dim_taxi_zonestable265 件の静的な zone。dbt 実行ごとに原子的なテーブル入れ替え (全件再構築) で全体を再構築する。部分更新はない
dim_payment_typetable6 件の静的な種別。dim_taxi_zones と同じ理由
dim_vendortable3 件の vendor。同じ理由

エンジンの設定

モデルENGINEバージョンカラム理由
fact_tripsReplacingMergeTree(updated_at)updated_atトリップは訂正されうる。挿入ごとに updated_at を now() に設定することで、RMT のバックグラウンド重複排除時に最新バージョンが勝つ。delete_insert が正しさを担う主経路であり、RMT はセーフティネット
agg_hourly_zone_tripsReplacingMergeTree(updated_at)updated_atローリング再計算は同じ (hour_bucket, zone_id) の組に対して集計を再挿入する。RMT により、バックグラウンドマージ時に陳腐な集計が取り除かれる
dim_taxi_zonesMergeTree()—dbt による全件リロードは、実行ごとに原子的なテーブル入れ替え (全件再構築) を意味する。重複は蓄積しえない。重複排除は不要
dim_payment_typeMergeTree()—dim_taxi_zones と同じ
dim_vendorMergeTree()—dim_taxi_zones と同じ

インクリメンタル戦略

モデルunique_keyincremental_strategyインクリメンタルフィルタこのフィルタを選ぶ理由
fact_tripstrip_iddelete_insertWHERE 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_insertWHERE pickup_at >= now() - INTERVAL 2 HOUR2 時間のローリングウィンドウにより境界の時間帯が必ず再集計され、途中までの時間帯のカウントが常に訂正されます。max(pickup_at) の高水位マークでは境界の時間帯を恒久的に過小カウントします

FINAL の配置

モデルFROM 句に FINAL を付けるか?理由
stg_trips付ける — FROM trips_raw FINALtrips_raw は ReplacingMergeTree であり、移行スクリプトのリトライやカットオーバー後のプロデューサーのリトライにより trip_id が重複した行を持ちえます。stg_trips が唯一の強制ポイントです: ここで重複排除し、下流のすべてのモデル (int_trips_enriched、fact_trips、agg_hourly_zone_trips) がクリーンなデータを受け取るようにします
int_trips_enriched付けないRMT テーブルではなく stg_trips (ビュー) から読みます。ビューに FINAL は無関係です
fact_trips付けない (モデル本体では)delete_insert により、実行が完了するたびに fact_trips はクリーンに保たれます。モデル内に FINAL を加えると、{{ this }} から max(updated_at) を読む is_incremental() のサブクエリにも無駄に適用されてしまいます。ダッシュボードと dbt テストは、fact_trips を直接クエリする際に外側で FINAL を使います

これは完成した例です。あなたの migration-plan.md は、ここでの主要な判断と一致しているはずです — あるいは、別の選択をした理由を明示的に記述してください。

このページの内容

JA