ワークト例: 完成した計画
NYC タクシーのワークロードについて記入済みの移行計画。自分で書き上げたあとに比較するためのもの。
これは 5 つのワークシートすべて
(1、
2、
3、
4、
5) を NYC
Taxi のワークロードに適用した完全な解答です。次の用途に使ってください:
- 各セクションを終えたあと、自分のワークシートの答えを確認する
- Part 3 が実装している判断の背後にある理由を理解する
- 自分が別の選択をした場合に、Part 3 の Decision Alignment 表と比較する
これは解答例です — 自分の計画としてここに記入しないでください。代わりに migration-plan.md に記入してください。
| メトリクス | 値 |
|---|
| テーブル総数 | 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) |
| Stream | 1 (TRIPS_RAW 上の TRIPS_CDC_STREAM) |
| Task | 2 (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) |
| オブジェクト | 種別 | スキーマ | 行数 | 複雑度グレード | 備考 |
|---|
trips_raw | テーブル | raw | 約 50M | B | _synced_at をバージョンカラムに持つ RMT。バルクと CDC の重複により重複排除が必要。stg_trips は FINAL を使う必要がある |
stg_trips | dbt ビュー | staging | — | B | TRIP_METADATA に対する JSONExtract。JSON パスのテストが必要 |
stg_taxi_zones | dbt ビュー | staging | — | A | パススルー。ごく簡単 |
int_trips_enriched | dbt Ephemeral | staging | — | A | CTE。SQL の差異は親モデル側で処理される |
fact_trips | dbt Incremental | analytics | 約 50M | C | RMT エンジン。delete_insert。QUALIFY の書き換え。FINAL が必要 |
agg_hourly_zone_trips | dbt Incremental | analytics | 約 140K | B | RMT。2 時間のローリング再計算ウィンドウ。パーティション境界を慎重にテストすること |
dim_taxi_zones | dbt テーブル | analytics | 265 | A | 静的な参照データ。全件リロード。ごく簡単 |
dim_payment_type | dbt テーブル | analytics | 6 | A | 静的な参照データ。ごく簡単 |
dim_vendor | dbt テーブル | analytics | 3 | A | 静的な参照データ。ごく簡単 |
taxi_zones_dict | ディクショナリ | analytics | 265 | B | ClickHouse 固有の構文。クエリ時に dictGet() |
mv_hourly_revenue | リフレッシュ可能な MV | analytics | — | B | REFRESH EVERY 構文。原子的な置き換えを検証すること |
TRIPS_CDC_STREAM / CDC_CONSUME_TASK | Snowflake Stream + Task | — | — | D | ClickHouse に相当機能なし。Part 3 ではプロデューサーの直接カットオーバーで置き換える |
| テーブル | エンジン | バージョンカラム | 理由 |
|---|
trips_raw | ReplacingMergeTree(_synced_at) | _synced_at | Python の移行スクリプト (scripts/02_migrate_trips.py) はバッチをリトライして同じ trip_id を再挿入する可能性があります。カットオーバー後も、ライブのプロデューサーが一時的な障害でリトライすることがあります。_synced_at DateTime DEFAULT now() は INSERT 時に自動設定されるため、後のリトライはより大きなタイムスタンプを持ち、RMT は最新の書き込みを残します。stg_trips は FINAL でクエリし、下流のモデルが動く前に重複排除を強制します。 |
fact_trips | ReplacingMergeTree(updated_at) | updated_at | トリップは訂正されえます (運賃の調整、ステータス変更)。同じ trip_id が更新後の値で再挿入されます。updated_at は訂正ごとに単調増加するため、RMT の重複排除では大きい値が勝ちます。常に FINAL でクエリしてください。 |
agg_hourly_zone_trips | ReplacingMergeTree(updated_at) | updated_at | dbt は直近 2 時間を再計算して再挿入します。RMT なしでは古い集計と新しい集計が蓄積し、二重カウントになります。dbt 実行ごとに updated_at を now() に設定することで、最新の値が勝つようにします。 |
dim_taxi_zones | MergeTree() | — | dbt による全件リロード (原子的なテーブル入れ替え (全件再構築))。重複が蓄積することはありません。重複排除は不要です。 |
dim_payment_type | MergeTree() | — | 同様 — 全件リロード。 |
dim_vendor | MergeTree() | — | 同様 — 全件リロード。 |
mv_hourly_revenue | MergeTree() | — | REFRESHABLE MV は REFRESH ごとに結果セット全体を原子的に置き換えます。upsert はありません。 |
| テーブル | 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 は慣習的で、可読性にも役立ちます。 |
| カラム | Snowflake の型 | ClickHouse の型 | 判断の根拠 |
|---|
TRIP_METADATA | VARIANT | String | 生の JSON をそのまま保持します。JSONExtract* はクエリ時に任意のパスを扱えます。Map(String,String) ではネスト構造が失われ、Tuple は固定スキーマを要求します。任意の JSON には String が安全な選択です。 |
PICKUP_DATETIME / PICKUP_AT | TIMESTAMP_NTZ(9) | DateTime64(3, 'UTC') | トリップのタイムスタンプにはミリ秒精度で十分です。ナノ秒 (9) は過剰です。'UTC' はタイムゾーンを明示し、時間範囲の集計での DST に起因する意外な挙動を避けます。 |
PICKUP_LOCATION_ID | INTEGER | UInt16 | 値は 1–265。UInt8 の最大値は 255 (小さすぎます)。UInt16 の最大値は 65535 (適切)。Int32 の 4 バイトに対し 2 バイト — 5,000 万行、1 カラムあたり非圧縮で約 95MB 節約します。 |
VENDOR_ID | INTEGER | UInt8 | 値は 1–3。UInt8 の最大値は 255 — 適切です。1 行あたり 1 バイト。 |
DRIVER_RATING | FLOAT | Nullable(Float32) | しばしば NULL です (すべてのトリップに評価があるわけではありません)。Nullable は正しい null セマンティクスを保ちます。1.0–5.0 の範囲には Float32 で十分です。Float64 は意味のある精度を加えずにストレージを浪費します。 |
UPDATED_AT | TIMESTAMP_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::FLOAT | JSONExtractFloat(trip_metadata, 'driver', 'rating') |
TRIP_METADATA:surge_multiplier::FLOAT | JSONExtractFloat(trip_metadata, 'surge_multiplier') |
QUALIFY ROW_NUMBER() OVER (PARTITION BY pickup_location_id ORDER BY fare_amount DESC) <= 10 | SELECT ... FROM (SELECT ..., ROW_NUMBER() OVER (...) AS rn FROM ...) WHERE rn <= 10 |
MERGE INTO fact_trips ... WHEN MATCHED THEN UPDATE | dbt の delete_insert インクリメンタル — キーが一致する行を DELETE し、その後すべての新しい行を INSERT |
| ウェーブ | オブジェクト | 依存関係 | 備考 |
|---|
| ウェーブ 0 | trips_raw (スキーマ)、dim_taxi_zones、dim_payment_type、dim_vendor | なし | dbt が空のテーブルを作成します。ディメンションテーブルは静的な参照データから直ちに投入されます (トリップへの依存なし)。実行: dbt run --select trips_raw dim_* |
| ウェーブ 1 | Python によるバルクロード (scripts/02_migrate_trips.py) | ウェーブ 0 (trips_raw のスキーマが存在していること) | Snowflake の TRIPS_RAW から 5,000 万行。--resume で再開可能。scripts/01_verify_migration.sh で行数を検証。約 40〜50 分。 |
| ウェーブ 2 | stg_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 がその上に構築されます。 |
| ウェーブ 3 | taxi_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 を投入します。 |
| オブジェクト | リスク | 検証方法 |
|---|
fact_trips | FINAL なしのクエリはマージ遅延の間に過大カウントします。対象外のパーティションを削除しないよう、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_trips | 2 時間のローリング再計算ウィンドウは削除範囲を正しく限定する必要があります。広すぎると古い集計が削除され、狭すぎると陳腐な集計が残ります。 | 特定の (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) が重複挿入をべき等に処理します。 |
データ移動: 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 コミュニティの標準的な推奨です。
| 基準 | しきい値 | 測定方法 |
|---|
| 行数の一致 | 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 の出力を手動で比較 |
| モデル | マテリアライゼーション | 理由 |
|---|
stg_trips | view | trips_raw を読み取ってクリーンにする。このモデル自体に更新はない。ストレージコストがゼロ。常に現在のソースの状態を反映する |
stg_taxi_zones | view | 同様 — ソーステーブルのパススルーなクリーンアップ |
int_trips_enriched | ephemeral | fact_trips からのみ使われる純粋な結合ロジック。CTE としてインライン展開することで冗長な物理テーブルを避ける。直接クエリするモデルはない |
fact_trips | incremental | トリップは事後に訂正されうる。実行ごとに新規行と更新行のみを処理すべき |
agg_hourly_zone_trips | incremental | 2 時間のローリング再計算はインクリメンタルなパターン — 5,000 万行すべてではなく最近の行を処理する |
dim_taxi_zones | table | 265 件の静的な zone。dbt 実行ごとに原子的なテーブル入れ替え (全件再構築) で全体を再構築する。部分更新はない |
dim_payment_type | table | 6 件の静的な種別。dim_taxi_zones と同じ理由 |
dim_vendor | table | 3 件の vendor。同じ理由 |
| モデル | ENGINE | バージョンカラム | 理由 |
|---|
fact_trips | ReplacingMergeTree(updated_at) | updated_at | トリップは訂正されうる。挿入ごとに updated_at を now() に設定することで、RMT のバックグラウンド重複排除時に最新バージョンが勝つ。delete_insert が正しさを担う主経路であり、RMT はセーフティネット |
agg_hourly_zone_trips | ReplacingMergeTree(updated_at) | updated_at | ローリング再計算は同じ (hour_bucket, zone_id) の組に対して集計を再挿入する。RMT により、バックグラウンドマージ時に陳腐な集計が取り除かれる |
dim_taxi_zones | MergeTree() | — | dbt による全件リロードは、実行ごとに原子的なテーブル入れ替え (全件再構築) を意味する。重複は蓄積しえない。重複排除は不要 |
dim_payment_type | MergeTree() | — | dim_taxi_zones と同じ |
dim_vendor | MergeTree() | — | dim_taxi_zones と同じ |
| モデル | unique_key | incremental_strategy | インクリメンタルフィルタ | このフィルタを選ぶ理由 |
|---|
fact_trips | trip_id | delete_insert | WHERE 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_insert | WHERE pickup_at >= now() - INTERVAL 2 HOUR | 2 時間のローリングウィンドウにより境界の時間帯が必ず再集計され、途中までの時間帯のカウントが常に訂正されます。max(pickup_at) の高水位マークでは境界の時間帯を恒久的に過小カウントします |
| モデル | FROM 句に FINAL を付けるか? | 理由 |
|---|
stg_trips | 付ける — FROM trips_raw FINAL | trips_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 は、ここでの主要な判断と一致しているはずです — あるいは、別の選択をした理由を明示的に記述してください。