Snowflake MigrationClickHouse Workshops

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_tripsstagingView手順1(毎回の実行)型キャスト、trip_metadata への JSONExtract*
stg_taxi_zonesstagingView手順1(毎回の実行)zone ディメンションのパススルー
int_trips_enrichedstagingEphemeral—(CTE としてインライン化)すべてのディメンション結合。物理テーブルなし
fact_tripsanalyticsIncremental手順1trip_id をキーにした delete_insert、updated_at でウォーターマーク。ORDER BY (toStartOfMonth(pickup_at), pickup_at, trip_id)
agg_hourly_zone_tripsanalyticsIncrementalモジュール05、カットオーバー後ローリング2時間ウィンドウ。増分フィルターは稼働中プロデューサーの行にのみ一致 — 下の手順1を参照
dim_taxi_zonesanalyticsTable手順1実行ごとにフルリロード。手順2の zone ディクショナリのソース
dim_payment_typeanalyticsTable手順1実行ごとにフルリロード
dim_vendoranalyticsTable手順1実行ごとにフルリロード
dim_dateanalyticsTable手順12009〜2029年の静的な日付スパイン。実行ごとにフルリロード
mv_live_trip_feedanalyticsMaterialized 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 run

dbt は 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 million
SELECT 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: 265
cd "$(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パスの手順で、そのギャップを意図的に 閉じます。いまプロデューサーを止めないでください。

このページの内容

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