Snowflake MigrationClickHouse Workshops

MergeTree エンジン

MergeTree ファミリーからエンジンを選び、働きに見合う ORDER BY キーを設計する。

ClickHouse は、すべてのデータを MergeTree エンジンのいずれかのバリアントに支えられたテーブルに保存します。Snowflake から来た人にとって、これに相当する概念はありません — Snowflake はストレージに関する判断をすべて内部で処理します。ClickHouse では、適切なエンジンを選ぶのは自分の責任であり、間違えると黙って不正な結果が生まれます。

このガイドでは、NYC タクシーのラボで使うエンジンと、Snowflake からの移行者が必ずつまずく落とし穴を扱います。


MergeTree とは

MergeTree は ClickHouse の主要なストレージエンジンです。データは パート と呼ばれるイミュータブルなカラムナファイルに書き込まれます。ClickHouse はバックグラウンドで定期的にパートをマージし、エンジンのルールに従ってソート、圧縮、そして必要に応じて変換を行います。

重要な帰結はこうです。マージが起きるまで、読み取りは1つの行の複数バージョンを見る可能性があります。 ほとんどのエンジンはこれを透過的に扱いますが、一部(特に ReplacingMergeTree)では、正しいクエリを書くためにマージのライフサイクルを理解する必要があります。

MergeTree テーブルを作成するときは ORDER BY を必ず指定します。これが決めるのは次のとおりです。

  1. 各パート内でのデータの物理的なソート順
  2. プライマリインデックス(スパースかつブロック単位で、メモリ上に保持される)
  3. 重複排除を行うエンジンでは、重複排除の「キー」を定義する列

primary key、クラスタ化インデックス、分散キーといった別個の概念はありません。ORDER BY がそれらすべてを一度に兼ねています。


MergeTree

使う場面: テーブルが挿入のみ、または更新が外部で処理される場合。重複排除は不要。

CREATE TABLE default.some_events (
    event_id      String,
    occurred_at   DateTime64(3, 'UTC'),
    payload       String
)
ENGINE = MergeTree()
ORDER BY (occurred_at, event_id);

特性:

  • 挿入は新しいパートとしてデータを追記する
  • 重複排除なし — 重複行はそのまま保持される
  • マージはストレージと圧縮を最適化するが、論理的な内容は変えない
  • クエリは ORDER BY のプレフィックス範囲に一致するすべてのパートを読む

うまくいかないとき: 同じ行を2回挿入すると(例: ネットワーク障害後のリトライ)、両方の行がクエリ結果に現れます。重複が起こり得ない、真に挿入のみのパイプラインではこれが正しい挙動です。CDC 更新やリトライ可能なロードを受けるテーブルには、ReplacingMergeTree を使ってください。


ReplacingMergeTree

使う場面: 行が更新され得る場合(例: 運賃の修正、ステータスの変更)。クエリ結果でキーごとに1行だけを得たい場合。

CREATE TABLE analytics.fact_trips (
    trip_id       String,
    pickup_at     DateTime64(3, 'UTC'),
    fare_amount   Float64,
    updated_at    DateTime64(3, 'UTC'),
    -- ...
)
ENGINE = ReplacingMergeTree(updated_at)
ORDER BY (toStartOfMonth(pickup_at), pickup_at, trip_id);

特性:

  • バックグラウンドのマージ中に、同じ ORDER BY キーを持つ行が重複排除される。残るのは バージョン列の値が最大の行 だけ
  • バージョン列(ここでは updated_at)がどの行を残すかを決める。値が大きい = より新しい = 残る
  • 重複排除は 非同期 — マージが走るまでは、古いバージョンと新しいバージョンが共存する

決定的な落とし穴: 重複排除の遅延

マージの合間には、FINAL を付けないクエリは行の すべての バージョンを見ます。

-- This may return multiple rows for the same trip_id
-- if the row has been updated since the last merge
SELECT * FROM analytics.fact_trips WHERE trip_id = 'abc123';

-- This returns exactly one row per trip_id, applying deduplication at query time
SELECT * FROM analytics.fact_trips FINAL WHERE trip_id = 'abc123';

FINAL は読み取り時に重複排除を強制します。重複キーがないかすべてのパートを確認しなければならないため、FINAL なしの読み取りより遅くなります。NYC タクシーのラボでは、fact_trips に対するすべてのクエリが FINAL を使います。

うまくいかないとき:

  • ポイントルックアップで FINAL を省く → 黙って重複行を返す。集計が過大になる
  • 誤ったバージョン列を使う(更新時に増加しない列) → 古い値が勝つ
  • 可変なテーブルに RMT ではなく MergeTree を使う → すべてのバージョンが蓄積し、行数が無制限に増える
  • 重複排除が同期的だと期待する → ETL ジョブが挿入直後に読み取り、重複を見てしまう

RMT と dbt: delete_insert インクリメンタル戦略は、挿入する前に取り込むバッチのキー範囲にある行を削除するため、テーブルにはそもそも重複が生じません。安全のため FINAL は依然として推奨されますが、dbt の戦略が正しければ重要度は下がります。


AggregatingMergeTree

使う場面: テーブルが部分集計の状態を保持し、それをバックグラウンドのマージ中にマージし、クエリ時に結合する場合。

CREATE TABLE analytics.agg_hourly_revenue (
    hour_bucket   DateTime,
    borough       String,
    fare_sum      AggregateFunction(sum, Float64),
    trip_count    AggregateFunction(count, UInt64)
)
ENGINE = AggregatingMergeTree()
ORDER BY (hour_bucket, borough);

特性:

  • 同じ ORDER BY キーを持つ行が、集計関数のコンバイナのロジックによってマージされる
  • クエリ時には -Merge サフィックスのコンバイナを使う: sumMerge(fare_sum)、countMerge(trip_count)
  • 通常は、生の挿入を部分状態へ変換する Materialized View から供給される

どんなときに使うか: AggregatingMergeTree は、部分状態が結合可能でなければならない事前集計データのためのものです。NYC タクシーのラボでは、agg_hourly_zone_trips は実行ごとに dbt が再構築します — 部分状態を積み上げるテーブルではなく、全置換のテーブルです。ここでは ReplacingMergeTree を使ってください。

うまくいかないとき: クエリ時に sumMerge(fare_sum) ではなく sum(fare_sum) を使うと、バイナリの集計状態を Float64 として扱い、ゴミのような数値を返します。これは黙って起こる正しさのエラーです。


CollapsingMergeTree

使う場面: 「sign 行」を挿入して行を削除・更新する必要がある場合(挿入は sign=1、取り消しは sign=-1)。あまり一般的ではありませんが、イベントベースの CDC パターンでは有用です。

ENGINE = CollapsingMergeTree(sign)

マージ中に、同じキーで sign=1 と sign=-1 の行のペアが互いに打ち消し合います。NYC タクシーのラボでは使いません — このワークロードの挿入リトライのパターンには、バージョン列を持つ ReplacingMergeTree のほうが単純です。


TTL 付きの MergeTree

どの MergeTree バリアントにも、時間ベースのデータ失効を追加できます。

CREATE TABLE default.trips_raw (
    trip_id    String,
    pickup_at  DateTime64(3, 'UTC'),
    _synced_at DateTime DEFAULT now(),
    -- ...
)
ENGINE = ReplacingMergeTree(_synced_at)
ORDER BY (pickup_at, trip_id)
TTL toDate(pickup_at) + INTERVAL 2 YEAR;

TTL はバックグラウンドのマージ中に作動します。期限切れの行は、パートがマージされるときに取り除かれます。このラボでは TTL を設定しません — 4年分のデータをすべて保持します。本番では、ストレージコストを管理するうえで TTL は不可欠です。


エンジンの選び方: 判断のツリー

Does the table receive UPDATE or DELETE operations?
├── No (insert-only, e.g., event log, append-only stream)
│   └── MergeTree()
└── Yes
    ├── Do rows have a version/timestamp column that increases on update?
    │   ├── Yes → ReplacingMergeTree(version_col)
    │   └── No (full reload, e.g., dim tables rebuilt by dbt)
    │       └── MergeTree() — dbt atomic table swap (full rebuild) handles "upsert"
    └── Is the table a pre-aggregated accumulator with combinable states?
        └── AggregatingMergeTree()

NYC タクシーのラボでは:

テーブルエンジン理由
trips_rawReplacingMergeTree(_synced_at)マイグレーションスクリプトのリトライやカットオーバー後のプロデューサーのリトライで、同じ trip_id が2回書かれ得る。_synced_at DEFAULT now() により、あとの書き込みが勝つ
fact_tripsReplacingMergeTree(updated_at)trip は修正され得る。バージョンは updated_at
agg_hourly_zone_tripsReplacingMergeTree(updated_at)ローリングでの再計算 = upsert。バージョンは updated_at
dim_* テーブルMergeTreedbt による全件リロード。部分更新はない
mv_hourly_revenueRefreshable MVスケジュールで実行され、毎回結果全体を置き換える

ORDER BY の設計

ORDER BY は ClickHouse のテーブルにおいて最も重要な性能上の判断です。これが決めるのは次のとおりです。

  1. プライマリインデックスの効率 — ORDER BY のプレフィックス列で絞り込むクエリは、無関係なブロックを読み飛ばせる
  2. 圧縮率 — ソート済みのデータはよく圧縮される(近い値が隣接するため)
  3. 重複排除のキー(RMT/AMT の場合) — 2つの行が重複とみなされるのは、ORDER BY の列が一致する場合のみ

ORDER BY を設計するルール:

  1. カーディナリティの低い列を先に 置く(例: borough、payment_type): 同じ値を共有する行が多くなるため、インデックスがより多くのブロックを読み飛ばせる
  2. カーディナリティの高い列を最後に 置く(例: trip_id、UUID): 範囲は絞れるが、先頭に置くと圧縮が効きにくい
  3. 列はソーススキーマからではなく、実際のクエリの絞り込み条件 から導く
  4. RMT のテーブルでは、最後の列を 行を一意に識別する列 にする(ビジネスキーごとに1行を保証する)

アンチパターン: ソースの primary key をそのまま ORDER BY にコピーすること。Snowflake の TRIPS_RAW に明示的なソートがない場合、Snowflake のスキーマ順(trip_id を先頭)をコピーすると ClickHouse ではランダムな ORDER BY になり、どの分析クエリでもブロックの読み飛ばしが効きません。

fact_trips での導出例:

Q1〜Q7 のクエリはすべて、何らかの形で pickup_at を絞り込んでいます。

  • Q1: WHERE pickup_at >= ...
  • Q2: ORDER BY week, pickup_location_id
  • Q3: WHERE pickup_at >= CURRENT_DATE - 7
  • Q4: GROUP BY DATE_TRUNC('day', pickup_at)

したがって pickup_at は ORDER BY に含めるべきで、しかも先頭寄りに置くべきです。最初の列に toStartOfMonth(pickup_at) を使うと、より粗い粒度のプレフィックスができ、PARTITION BY 句がなくてもパーティションレベルの枝刈りに相当する効果が得られます。trip_id は RMT の一意性のために最後に置きます。

結果: ORDER BY (toStartOfMonth(pickup_at), pickup_at, trip_id)


PARTITION BY

PARTITION BY は任意で、ORDER BY とは別物です。物理的なディレクトリのパーティションを作り、各パーティションが独立したパートの集合になります。

PARTITION BY toYYYYMM(pickup_at)

PARTITION BY を使う場面:

  • ある期間全体を効率よく DROP する必要がある(ALTER TABLE DROP PARTITION '202401')
  • TTL を行単位ではなく月単位で働かせたい
  • テーブルが非常に大きく(1TB 超)、パーティション単位のメタデータがクエリプランニングに役立つ

ORDER BY の代わりに PARTITION BY を使ってはいけません。 よくある間違いは、toYYYYMM(date) を PARTITION BY に入れて ORDER BY から省くことです — これはパーティション内でのブロック単位の読み飛ばしを妨げます。

NYC タクシーのラボでは PARTITION BY は不要です — データセットは5,000万行(圧縮後で約8GB)で、単一パーティションの性能範囲に十分収まっています。


主な落とし穴のまとめ

落とし穴帰結対処
可変なデータに誤ったエンジン重複行が黙って蓄積するReplacingMergeTree + FINAL を使う
RMT のクエリで FINAL がないマージ遅延の間、集計が過大になるRMT テーブルに対するすべての分析クエリに FINAL を付ける
ソーススキーマから写した ORDER BYクエリが遅い。ブロックの読み飛ばしが効かないORDER BY は実際のクエリの絞り込み条件から導く
ORDER BY の先頭がカーディナリティの高い列インデックスの選択性が低いカーディナリティは低い順に、高い列は最後に
AggregateFunction 列を sumMerge() ではなく sum() でクエリ黙って数値がゴミになるAggregatingMergeTree には必ず -Merge コンバイナを使う
RMT のバージョン列が単調増加しない古いバージョンがランダムに勝つ更新時に必ず now() が設定されるタイムスタンプを使う

このページの内容

JA