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 つのマテリアライゼーションをサポートします。好みではなく、更新パターンに基づいて選んでください。
| マテリアライゼーション | 物理オブジェクト | 使いどころ |
|---|---|---|
view | ClickHouse ビュー | staging モデル: ソースデータのクリーンアップと型キャスト。ストレージコストなし。クエリごとに再評価される |
ephemeral | オブジェクトなし (CTE としてインライン展開) | 複数の staging モデルを JOIN で組み合わせる intermediate モデル。冗長な物理テーブルの作成を避ける |
table | staging リレーション内に完全な置き換えを構築し、EXCHANGE TABLES (古いバージョンではリネームのペア) で原子的に入れ替える。入れ替え後に古いテーブルを DROP | dbt run ごとに全件置き換えられる小さなディメンションテーブル。部分更新が不要な場合。注意: 大きなテーブルで全件再構築は現実的ではありません — 数千行を超えるテーブルには incremental を使ってください。 |
incremental | 初回実行では CREATE TABLE、以降は選択的な UPDATE パターン | 実行ごとに新規/変更行のみを処理すべきファクトテーブルおよび事前集計テーブル |
materialized_view | ClickHouse 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_by | primary key / ソート順 | CREATE TABLE の ORDER BY ...。省略時は tuple() |
+unique_key | delete_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を設定します。
何をするのか
インクリメンタル実行のたびに:
- 到着したバッチのいずれかの行と
unique_keyが一致する行を、ターゲットテーブルから DELETE する - 到着したバッチのすべての行を 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 %}