Snowflake MigrationClickHouse Workshops

ClickHouse 上の dbt

dbt-clickhouse の設定: delete_insert インクリメンタル戦略、ReplacingMergeTree モデル、リフレッシュ可能な materialized view。

このガイドでは、Part 3 で使う dbt-clickhouse 固有のパターンを扱います。ワークシート 1–4 を終えたあと、ワークシート 5 (dbt モデル設計) の前に読んでください。

dbt-snowflake から来た場合、dbt の概念のほとんどは同一です — source、ref、テスト、マクロ、staging/intermediate/analytics のレイヤーパターン。変わるのは ClickHouse 固有の設定レイヤー、すなわち engine、order_by、インクリメンタル戦略、そして FINAL のセマンティクスです。


1. マテリアライゼーションの種類

dbt-clickhouse は 5 つのマテリアライゼーションをサポートします。好みではなく、更新パターンに基づいて選んでください。

マテリアライゼーション物理オブジェクト使いどころ
viewClickHouse ビューstaging モデル: ソースデータのクリーンアップと型キャスト。ストレージコストなし。クエリごとに再評価される
ephemeralオブジェクトなし (CTE としてインライン展開)複数の staging モデルを JOIN で組み合わせる intermediate モデル。冗長な物理テーブルの作成を避ける
tablestaging リレーション内に完全な置き換えを構築し、EXCHANGE TABLES (古いバージョンではリネームのペア) で原子的に入れ替える。入れ替え後に古いテーブルを DROPdbt run ごとに全件置き換えられる小さなディメンションテーブル。部分更新が不要な場合。注意: 大きなテーブルで全件再構築は現実的ではありません — 数千行を超えるテーブルには incremental を使ってください。
incremental初回実行では CREATE TABLE、以降は選択的な UPDATE パターン実行ごとに新規/変更行のみを処理すべきファクトテーブルおよび事前集計テーブル
materialized_viewClickHouse Materialized View自動リフレッシュされる集計。dbt の incremental とは異なります。標準の (トリガーベースの) MV は INSERT ごとに 1 回起動し、そのバッチしか見えません — 全期間の集計は計算できません。一方 REFRESHABLE MV はスケジュールに従ってクエリ全体を再実行するため、それが可能です。

Snowflake との主な違い: dbt-snowflake はストレージの詳細を内部で処理します。dbt-clickhouse では、table と incremental モデルに明示的な +engine 設定が必要です — dbt はこれを使って CREATE TABLE ... ENGINE = ... の DDL を生成します。

ビューに engine はありません。 うっかり view マテリアライゼーションに +engine を付けても、dbt-clickhouse はそれを無視します。engine が必要な永続ストレージを作るのは table と incremental のマテリアライゼーションだけです。

リフレッシュ可能な materialized view。 dbt-clickhouse の materialized_view マテリアライゼーションは refreshable の設定ブロック — interval (および任意で randomize) — を受け取り、生成する CREATE MATERIALIZED VIEW 文に REFRESH 句を直接出力します。このラボの mv_live_trip_feed モデルは refreshable を設定していないため、構築される MV にリフレッシュスケジュールはありません。


2. ClickHouse の設定を dbt で表現する

ClickHouse 固有の設定は、dbt のモデル設定として表現します。プロジェクト全体のデフォルトは dbt_project.yml、モデル固有の上書きはモデルの config() ブロックに記述します。

dbt_project.yml での記述

models:
  your_project:
    analytics:
      +schema: analytics
      +materialized: table
      +engine: "MergeTree()"          # default for all analytics tables

      fact_trips:
        +materialized: incremental
        +engine: "ReplacingMergeTree(updated_at)"   # overrides the default
        +incremental_strategy: delete_insert
        +unique_key: trip_id
        +order_by: "(toStartOfMonth(pickup_at), pickup_at, trip_id)"

モデルの config() ブロックでの記述

{{ config(
    materialized         = 'incremental',
    engine               = 'ReplacingMergeTree(updated_at)',
    incremental_strategy = 'delete_insert',
    unique_key           = 'trip_id',
    order_by             = '(toStartOfMonth(pickup_at), pickup_at, trip_id)'
) }}

どちらの方法も等価です。プロジェクト全体のパターンには dbt_project.yml が、モデル固有の上書きや設定を SQL と同じ場所に置きたい場合には config() ブロックが適しています。

主要な設定パラメータ

パラメータ何を制御するかClickHouse への対応
+engineテーブルのストレージエンジンCREATE TABLE の ENGINE = ...
+order_byprimary key / ソート順CREATE TABLE の ORDER BY ...。省略時は tuple()
+unique_keydelete_insert の重複排除キー挿入前にどの行を削除するかを決める
+incremental_strategyインクリメンタル実行がデータをどう更新するかClickHouse では delete_insert に設定

スコープのルール: dbt_project.yml の設定は親から子へカスケードします。モデルレベルの config() ブロックは常にプロジェクト設定に勝ちます。最も一般的なエンジンをプロジェクトのデフォルトに設定し、異なるモデルだけを上書きしてください。


3. delete_insert の仕組み

delete_insert は dbt-clickhouse コミュニティの標準的なインクリメンタル戦略です。Snowflake の MERGE INTO に最も近い相当物ですが、仕組みは異なります。

バージョン要件: delete_insert は ClickHouse の軽量削除 (lightweight delete) を使います。これは 22.8 で導入され (実験的)、23.3 以降で本番利用可能になりました。ClickHouse Cloud はこの要件を満たします。有効にするには、~/.dbt/profiles.yml の ClickHouse ターゲットに use_lw_deletes: true を追加するか、query_settings で allow_experimental_lightweight_delete=1 を設定します。

何をするのか

インクリメンタル実行のたびに:

  1. 到着したバッチのいずれかの行と unique_key が一致する行を、ターゲットテーブルから DELETE する
  2. 到着したバッチのすべての行を INSERT する
-- Step 1: dbt generates this DELETE
ALTER TABLE analytics.fact_trips
DELETE WHERE trip_id IN (SELECT trip_id FROM incoming_batch);

-- Step 2: dbt generates this INSERT
INSERT INTO analytics.fact_trips
SELECT * FROM incoming_batch;

Snowflake の MERGE INTO との違い

Snowflake の merge 戦略は、行単位の WHEN MATCHED THEN UPDATE / WHEN NOT MATCHED THEN INSERT を生成します。ClickHouse に MERGE INTO 文はありません。delete_insert は、バッチ削除に続く全件挿入によって同じ最終結果 — ユニークキーごとに 1 行 — を実現します。

ReplacingMergeTree との関係

delete_insert が正しさを担う主経路です。ReplacingMergeTree はセーフティネットです。

delete_insert の実行が正常に完了した場合: テーブルはクリーン (trip_id ごとに 1 行) で、重複はありません。

delete_insert の実行が途中で中断された場合 (DELETE 後、INSERT 前にクラッシュ): データは不正な状態になっている可能性が高く、削除された行が再挿入されていないかもしれません。次に成功する実行で正しい状態に復元されますが、失敗した DELETE とその再実行の間はテーブルにクエリしないでください。

何らかの理由で実行が重複を生じた場合: ReplacingMergeTree のバックグラウンドマージが最終的に重複を排除し、バージョンカラムの値が最大の行を残します。

delete_insert なしで RMT だけに頼ってはいけません — バックグラウンドマージは非同期であり、大きなテーブルでは数分から数時間かかることがあります。

append を使うとき

append は既存の行に触れずに新しい行を挿入します。行が決して更新されない、純粋な挿入専用テーブルには正しい戦略です — たとえばイミュータブルなイベントログや、ID の一意性が保証され訂正が入らない生の取り込みテーブルです。append にはバージョン要件も、ミューテーションのリスクもありません。

fact_trips にとって append は誤りです: トリップは事後に訂正されうるため (運賃の調整、ステータス変更)、同じ trip_id が新しい値で再度到着します。append では両方のバージョンが永続的に蓄積し、集計 (運賃の SUM、トリップの COUNT) は次のバックグラウンド RMT マージまで過大にカウントされます。行が更新されうる場合は常に delete_insert を使ってください。

なぜ merge 戦略ではないのか

merge 戦略 (delete_insert 以前のレガシーなデフォルト) は一時テーブルを作り、変更のない既存行と新しいバッチを入れて、元のテーブルを原子的に置き換えます。delete_insert と違って軽量削除を使わず、インクリメンタル実行ごとにテーブル全体を書き換えます。5,000 万行の fact_trips テーブルでは極端にコストが高くなります。delete_insert は現在のバッチにある行だけを処理しますが、merge はテーブルのすべての行に触れます。delete_insert を使ってください。


4. FINAL の配置戦略

ReplacingMergeTree の重複排除はバックグラウンドで行われます — ClickHouse はパートを非同期にマージします。マージの間、重複行は共存します。FINAL は読み取り時に同期的な重複排除を強制します。

dbt パイプラインにおける FINAL の置き場所

ReplacingMergeTree のソースから読み取り、クリーンな分析データを生成するレイヤーに置きます。

NYC Taxi のワークロードでは:

trips_raw (RMT)
    ↓
stg_trips (view): SELECT ... FROM trips_raw FINAL   ← FINAL goes here
    ↓
int_trips_enriched (ephemeral CTE)
    ↓
fact_trips (incremental, RMT)                        ← NO FINAL in model
    ↓
Dashboard queries: SELECT ... FROM fact_trips FINAL  ← FINAL goes here (externally)

stg_trips は trips_raw の重複排除を強制する唯一のポイントです。stg_trips を読む下流のモデルはすべて、自動的にクリーンで重複排除済みのソースデータを得ます。int_trips_enriched や fact_trips で FINAL は不要です。これらは (RMT テーブルではなくビューである) stg_trips から読むためです。

fact_trips を直接読むダッシュボードのクエリや dbt テストでは、外側で FINAL を使います。モデル自体に FINAL を埋め込まないのは、モデルのクエリ内のすべてのスキャン — {{ this }} から max(updated_at) を読む is_incremental() のサブクエリも含む — に適用されてしまうからです。

FINAL のパフォーマンスへの影響

FINAL は重複行の数に比例したレイテンシを加えます。よく整備された RMT テーブル (バックグラウンドマージが頻繁) では、解決すべき重複が少ないため FINAL のオーバーヘッドはごくわずかです。マージされていないパートが多い、ロードしたばかりのテーブルでは、FINAL は大幅に遅くなることがあります。

dbt テストと検証クエリでは、RMT テーブルに対して常に FINAL を使ってください。Snowflake とのレイテンシ比較が目的のベンチマーククエリについては、ClickHouse 側のクエリはすでに FINAL を使っています — したがって比較は公平です。


5. generate_schema_name マクロ

デフォルトでは、dbt はモデルのスキーマ名にプロファイルのターゲットスキーマ名を接頭辞として付けます。dbt プロファイルがスキーマ nyc_taxi_ch を対象にしている場合、+schema: analytics を持つモデルは analytics ではなく nyc_taxi_ch_analytics に作られます。

Snowflake では無害ですが (スキーマはデータベース内の名前空間)、スキーマそのものがデータベースである ClickHouse では扱いにくい名前になります。nyc_taxi_ch_analytics は有効な ClickHouse データベース名ですが、analytics より見苦しく、Part 3 の ClickHouse アーキテクチャで使うターゲットデータベース名とも一致しません。

対処は generate_schema_name マクロの上書きです:

-- macros/generate_schema_name.sql
{% macro generate_schema_name(custom_schema_name, node) -%}
  {%- if custom_schema_name is none -%}
    {{ target.schema | lower }}
  {%- else -%}
    {{ custom_schema_name | lower }}
  {%- endif -%}
{%- endmacro %}

このマクロは:

  • モデルが +schema: analytics を指定している場合、custom_schema_name をそのまま (小文字化して) 返す
  • カスタムスキーマを持たないモデルには、プロファイルのターゲットスキーマを (小文字化して) 返す

| lower フィルタは、ClickHouse の大文字小文字を区別する識別子ルールに合わせて、スキーマ名が一貫して小文字になることも保証します (Part 1 の Snowflake では | upper を使いました)。

置き場所: macros/generate_schema_name.sql — トップレベルの macros/ ディレクトリ内。dbt_project.yml で macro-paths: ["macros"] を設定しています。


まとめ: NYC Taxi の dbt 設定サマリ

# dbt_project.yml (abbreviated)
models:
  nyc_taxi_dbt_ch:
    staging:
      +schema: staging
      +materialized: view           # no engine — views need none

    intermediate:
      +schema: staging
      +materialized: ephemeral      # inlined as CTE

    analytics:
      +schema: analytics
      +materialized: table
      +engine: "MergeTree()"        # default for dim_* tables

      fact_trips:
        +materialized: incremental
        +engine: "ReplacingMergeTree(updated_at)"
        +incremental_strategy: delete_insert
        +unique_key: trip_id

      agg_hourly_zone_trips:
        +materialized: incremental
        +engine: "ReplacingMergeTree(updated_at)"
        +incremental_strategy: delete_insert
        +unique_key: [hour_bucket, zone_id]
-- stg_trips.sql (staging view — the FINAL enforcement point)
SELECT ... FROM {{ source('raw', 'trips_raw') }} FINAL

-- fact_trips.sql (incremental — no FINAL in model body)
SELECT ... FROM {{ ref('int_trips_enriched') }}
{% if is_incremental() %}
WHERE updated_at > (SELECT max(updated_at) FROM {{ this }})
{% endif %}

このページの内容

JA