04 dbt パイプラインの再構築
dbt-clickhouse で Medallion パイプラインを ClickHouse 上に再構築します — delete_insert の増分モデル、ReplacingMergeTree、リフレッシュ可能な materialized view — そして zone ディクショナリを作成します。
開始チェックポイント
モジュール03が完了していること。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 — モジュール
03の dbt run がこれら7つすべてをすでにビルドしています。analytics のテーブルはすべてまだ空ですが、
mv_live_trip_feed は例外で、そのビルド時のスナップショット1行をすでに保持しています。それ以外で
データがあるのは default.trips_raw のみで、およそ5,000万行です。Snowflake のプロデューサーは
まだ動いているので、ClickHouse はおおよそマイグレーションのウィンドウの長さだけ Snowflake より
遅れています。約30分を見込んでください。
なぜ必要か
モジュール03は、ClickHouse が5,000万行を保持できることを示しました。しかし、パイプライン が ClickHouse 上で動くことは示していません。staging view、増分の fact テーブル、ディメンションの リロード、壊れたモデルをパートナーが目にする前に捕まえるテスト、これらのことです。それがこの モジュールで再構築するものです。モジュール01と同じ Medallion のモデルを、dbt-snowflake の代わりに dbt-clickhouse で表現し、モジュール03がすでに作成したテーブルに対して実行します。
モデルのロジックは何も変わりません。stg_trips は依然として型キャストと JSON の抽出を行い、
int_trips_enriched は依然としてディメンションを結合し、fact_trips は依然として trip ごとに
1行に落ち着きます。変わるのはその下のマテリアライゼーションのレイヤーです。MERGE INTO は
ありません。Snowflake Task もありません。cluster_by もありません。このモジュールが示すのは、
マイグレーションが一度きりのデータダンプではないということです。パートナーのチームが毎日、同じ
スケジュールで、同じ dbt test の実行をゲートとして走らせているパイプラインは、その下の
ウェアハウスが変わっても動き続けます。
概念 — 内部の仕組み
このモデルセットは、Snowflake のパイプラインと4つの点で異なります。設定の完全なリファレンスは
dbt on ClickHouse です。以下は、手順1で
dbt run を実行する前に知っておくべき短縮版です。これらのモデルが置き換えるソース側の
パイプラインについては、
dbt on Snowflake を参照してください。
1. delete_insert が MERGE を置き換える。 ClickHouse には MERGE INTO 文がありません。
Snowflake のパイプラインが fact_trips と agg_hourly_zone_trips を upsert するために
incremental_strategy: merge を使っていた箇所では、ClickHouse のモデルは
incremental_strategy: delete_insert を使います。dbt は入ってくるバッチの unique_key に
一致する行を削除し、それからバッチを挿入します。fact_trips では unique_key は trip_id で、
増分のフィルターは pickup_at ではなく updated_at をウォーターマークにします。運賃の訂正は
同じ trip_id を同じ pickup_at で、しかしより新しい updated_at で再挿入するので、
pickup_at をウォーターマークにすると黙って取りこぼしてしまいます。
2. ReplacingMergeTree は delete_insert の下にある安全網であり、その代替ではない。
どちらの増分モデルも ReplacingMergeTree(updated_at) として宣言されています。delete_insert の
実行が正常に完了すれば、テーブルはすでにキーごとに1行になっており、エンジンが後片付けするものは
ありません。実行が途中で中断された場合 — 削除のあと、挿入の前にクラッシュした場合 — バックグラウンドの
マージが最終的に残った行を重複排除し、updated_at が最大のものを残します。delete_insert が
やるはずの重複排除の仕事を ReplacingMergeTree だけに任せてはいけません。バックグラウンドの
マージは非同期で、この規模のテーブルでは数分から数時間遅れることがあります。
3. リフレッシュ可能な materialized view がスケジュールされたタスクを置き換える。 Snowflake の
パイプラインは、ローリングの集計を最新に保つためにストアドプロシージャを実行するスケジュール
タスクを使っていました。ClickHouse の dbt プロジェクトは代わりに mv_live_trip_feed を
materialized = 'materialized_view' と engine = 'ReplacingMergeTree(refreshed_at)' で宣言し、
リフレッシュ可能な materialized view として dbt run がビルドします。同じ効果を得るための
Snowflake の CREATE TASK DDL 30行超と比べてみてください。このモジュールではリフレッシュ間隔を
有効にはしません。それを行うには
ALTER TABLE analytics.mv_live_trip_feed MODIFY REFRESH EVERY ... を手動で実行する必要があり、
ラボはそれをスクリプト化していません(理由はモジュール05を参照)。
4. mv_live_trip_feed には Snowflake 側の対応物がそもそもない。 既存モデルの変換ではなく、
マイグレーションによって追加される新しい能力です。標準の ClickHouse の materialized view は
INSERT ごとに1回発火し、そのバッチの行しか見ないので、総 trip 数や平均運賃のような通期の集計を
正しく計算できません。REFRESHABLE な materialized view は代わりにクエリ全体 — ここでは
SELECT ... FROM {{ ref('fact_trips') }} — をスケジュールに従って再実行するので、リフレッシュの
たびにテーブル全体を見ます。このワークショップの Snowflake 側には、この選択肢はありませんでした。
dbt が実際にビルドするモデルと、それぞれにデータが入るタイミングは次のとおりです。
| モデル | レイヤー | マテリアライゼーション | データ投入 | 備考 |
|---|---|---|---|---|
stg_trips | staging | View | 手順1(毎回の実行) | 型キャスト、trip_metadata への JSONExtract* |
stg_taxi_zones | staging | View | 手順1(毎回の実行) | zone ディメンションのパススルー |
int_trips_enriched | staging | Ephemeral | —(CTE としてインライン化) | すべてのディメンション結合。物理テーブルなし |
fact_trips | analytics | Incremental | 手順1 | trip_id をキーにした delete_insert、updated_at でウォーターマーク。ORDER BY (toStartOfMonth(pickup_at), pickup_at, trip_id) |
agg_hourly_zone_trips | analytics | Incremental | モジュール05、カットオーバー後 | ローリング2時間ウィンドウ。増分フィルターは稼働中プロデューサーの行にのみ一致 — 下の手順1を参照 |
dim_taxi_zones | analytics | Table | 手順1 | 実行ごとにフルリロード。手順2の zone ディクショナリのソース |
dim_payment_type | analytics | Table | 手順1 | 実行ごとにフルリロード |
dim_vendor | analytics | Table | 手順1 | 実行ごとにフルリロード |
dim_date | analytics | Table | 手順1 | 2009〜2029年の静的な日付スパイン。実行ごとにフルリロード |
mv_live_trip_feed | analytics | Materialized view(リフレッシュ可能) | モジュール03の dbt run。リフレッシュ間隔は一度も有効化されない | Snowflake 側の対応物なし — 上の項目4を参照 |
手順1 — analytics レイヤーにデータを投入する
モジュール00で構築した dbt-clickhouse の venv をアクティブにし、dbt run を2回目として実行します。
モジュール03では、スキーマを作るために空のテーブルに対して一度実行済みです。今回の実行には
実データが背後にあります。trips_raw にはいま5,000万行が入っています。
cd "$(git rev-parse --show-toplevel)/workshop_public/snowflake_migration_lab/03-migrate-to-clickhouse"
source .venv/bin/activate
source .env && source .clickhouse_state
cd "$(git rev-parse --show-toplevel)/workshop_public/snowflake_migration_lab/03-migrate-to-clickhouse/dbt/nyc_taxi_dbt_ch"
dbt rundbt は dbt/nyc_taxi_dbt_ch/profiles.yml.example に基づいて ~/.dbt/profiles.yml から ClickHouse の
接続情報を読み取ります。モジュール03が空のスキーマを作るときに使ったのと同じプロファイルです。
期待される結果: 約8〜12分(増分モデルが5,000万行を処理します)。
この実行のあと agg_hourly_zone_trips は空になります。これは想定どおりで、失敗ではありません。
その増分フィルターは WHERE pickup_at >= now() - INTERVAL 2 HOUR であり、稼働中のプロデューサーが
書き込んだ行にしか一致しません。いま移行した行はすべて過去のデータなので、たった今から測った
2時間のウィンドウには1行も入りません。このテーブルは、モジュール05のカットオーバーが ClickHouse の
プロデューサーを起動するまで空のままです。壊れたパイプラインとしてデバッグに時間を使わないで
ください。
続いてテストスイートを実行します。
dbt test期待される結果: すべてのテストが成功します。
検証:
SELECT formatReadableQuantity(count()) FROM analytics.fact_trips FINAL;
-- Expected: ~50 million
SELECT count() FROM analytics.agg_hourly_zone_trips;
-- Expected: 0 (normal — populated after cutover in module 05)
SELECT count() FROM analytics.dim_taxi_zones;
-- Expected: 265手順2 — zone ディクショナリを作成する
analytics.dim_taxi_zones にはこれで NYC TLC の265の zone すべてが入りました。そのテーブルを
バックエンドとするインメモリのディクショナリ analytics.taxi_zones_dict をビルドし、下流の
クエリが JOIN の代わりに dictGet() で zone の borough を引けるようにします。
ディクショナリが結合に対して得るもの。 ディクショナリは一度メモリに読み込まれ、そのまま
ホットな状態を保ちます。それに対する参照は、以降のクエリでは実質的にコストゼロです。
dim_taxi_zones に対する JOIN は、実行のたびにディメンションテーブルを再読み込みして再マッチ
します。このような小さくてほとんど変わらないリファレンステーブル — 265行で、dbt run ごとに
フルリロードされます — では、このトレードオフは一方的です。モジュール05のベンチマークの実行は
taxi_zones_dict を dictGet で直接クエリするため、この手順はそのモジュールにとって任意の
追加要素ではなく、必須の依存関係です。
接続情報を source し、それから clickhouse-client か HTTP API のどちらかでディクショナリの
DDL を読み込みます。手元で使える方を選んでください。
cd "$(git rev-parse --show-toplevel)/workshop_public/snowflake_migration_lab/03-migrate-to-clickhouse"
source .clickhouse_state
# Via clickhouse-client
clickhouse-client \
--host "${CLICKHOUSE_HOST}" \
--port 9440 \
--user default \
--password "${CLICKHOUSE_PASSWORD}" \
--secure \
--multiquery \
< scripts/04_create_dictionary.sql
# Or via HTTP API
curl "https://${CLICKHOUSE_HOST}:${CLICKHOUSE_PORT}/" \
--user "default:${CLICKHOUSE_PASSWORD}" \
--data-binary @scripts/04_create_dictionary.sql検証:
-- Should return 'Manhattan' for zone 42
SELECT dictGet('analytics.taxi_zones_dict', 'borough', toUInt16(42));
-- Should show status = LOADED, element_count = 265
SELECT name, status, element_count
FROM system.dictionaries
WHERE name = 'taxi_zones_dict';完了の確認方法
SELECT formatReadableQuantity(count()) FROM analytics.fact_trips FINAL;
-- Expected: ~50 millionSELECT count() FROM analytics.agg_hourly_zone_trips;
-- Expected: 0ここではゼロが正しく、失敗ではありません。agg_hourly_zone_trips の増分フィルター
(WHERE pickup_at >= now() - INTERVAL 2 HOUR)は稼働中のプロデューサーが書き込んだ行にしか
一致せず、いま ClickHouse にある行はすべてモジュール03のマイグレーションスクリプトが移した過去の
データです。now() から見て2時間以内のものは1行もありません。このテーブルは、モジュール05の
カットオーバーが ClickHouse のプロデューサーを起動したあとにだけ埋まります。それまでは、この
テーブルを使うダッシュボードのチャートにはデータが表示されませんが、それは想定どおりであり、
デバッグすべきものではありません。
SELECT count() FROM analytics.dim_taxi_zones;
-- Expected: 265cd "$(git rev-parse --show-toplevel)/workshop_public/snowflake_migration_lab/03-migrate-to-clickhouse/dbt/nyc_taxi_dbt_ch"
dbt test期待される結果: すべてのテストが成功します。
-- Should return a borough name, e.g. 'Manhattan'
SELECT dictGet('analytics.taxi_zones_dict', 'borough', toUInt16(42));終了状態
analytics レイヤーにデータが入り、テスト済みです。fact_trips はおよそ5,000万行を保持し、
dim_taxi_zones、dim_payment_type、dim_vendor、dim_date は完全に読み込まれ、dbt test は
端から端まで成功し、analytics.taxi_zones_dict は稼働して dictGet() 経由で borough を返します。
agg_hourly_zone_trips はまだ空です。これは不具合ではなく設計どおりで、モジュール05の
カットオーバーまでそのままです。
ダッシュボード、ClickHouse と Snowflake のベンチマーク、そしてカットオーバーは、このモジュールでは なくモジュール05です。
Snowflake のプロデューサーはまだ動いており、Snowflake と ClickHouse の間のギャップもまだ 開いたままです。 このモジュールではプロデューサーもマイグレーションスクリプトも触っていません。 モジュール05が、モジュール03で予告した同じ制御された2パスの手順で、そのギャップを意図的に 閉じます。いまプロデューサーを止めないでください。