Snowflake MigrationClickHouse Workshops

ClickHouse の運用

ターゲットサービスを動かす: クエリ実行、マージとパートの監視、ディクショナリ、そして Snowflake とは異なる運用上の習慣。

このドキュメントは、この移行ラボで遭遇する ClickHouse の概念を説明します — それが何であり、なぜ存在し、パート1で使った Snowflake の構成要素とどう違うのかです。


1. テーブルエンジン

ClickHouse は単一エンジンのデータベースではありません。作成するすべてのテーブルは エンジン を宣言しなければならず、それがディスク上でのデータの保存方法、重複の扱い、利用できる機能を決めます。エンジンの選択を誤ることは、ClickHouse のスキーマ設計で最もよくある間違いです。

MergeTree

ほぼすべての本番テーブルの基本となるエンジンです。

CREATE TABLE analytics.dim_taxi_zones (
    zone_id      UInt16,
    borough      String,
    service_zone String
) ENGINE = MergeTree()
ORDER BY zone_id;

何をするか: ClickHouse はデータを パート — ディスク上のソート済みで圧縮されたチャンク — に保存します。データを挿入すると、新しいパートが書き込まれます。バックグラウンドでは、ClickHouse が小さなパートをより大きなパートへ継続的に マージ し、ORDER BY キーによるソートを保ちます。エンジン名はここから来ています。

使う場面: 重複排除が不要で、挿入が追記のみまたはバルクロードであるテーブル全般(ディメンションテーブル、生イベントテーブル、ログテーブル)。

重要な性質: primary key の強制はありません。ORDER BY の値が同一の2行は、両方とも保存されます。重複排除が必要なら ReplacingMergeTree を使ってください。

ReplacingMergeTree(version_col)

重複排除のためのエンジンです。MergeTree に次のルールを加えます: バックグラウンドのマージ中に、2つの行が同じ ORDER BY キーを共有していたら、version_col の値が最大のものだけを残す。

CREATE TABLE analytics.fact_trips (
    trip_id     String,
    pickup_at   DateTime,
    total_amount Float64,
    updated_at  DateTime
) ENGINE = ReplacingMergeTree(updated_at)
ORDER BY (pickup_at, trip_id);

重要 — 結果整合性: 重複排除はバックグラウンドのマージ中にしか起きません。ある瞬間を切り取れば、テーブルには重複行が含まれている可能性があります。これを 結果整合性 と呼びます。クエリ時に完全に重複排除された結果を得るには、SELECT に FINAL を付けてください。

-- Without FINAL: may return duplicates if merges haven't run
SELECT * FROM analytics.fact_trips WHERE trip_id = 'abc';

-- With FINAL: forces deduplication at query time (slower, always correct)
SELECT * FROM analytics.fact_trips FINAL WHERE trip_id = 'abc';

dbt での使われ方: dbt-clickhouse アダプターは主要な仕組みとして delete_insert インクリメンタル戦略を使います — キーが一致する行を明示的に削除して新しい行を挿入するため、常に正しい結果になります。ReplacingMergeTree は、すり抜けた重複(たとえば部分的に失敗した挿入から生じたもの)を片付ける セーフティネット として働きます。

Snowflake での等価物: 直接の等価物はありません。Snowflake では MERGE INTO ... WHEN MATCHED THEN UPDATE を使いました。ClickHouse に MERGE 文はなく — ReplacingMergeTree と FINAL の組み合わせが同じ論理的結果を実現します。

Refreshable Materialized View

ClickHouse は2種類の materialized view をサポートします。

トリガー型 MV(従来型): INSERT ごとに実行され、新しく挿入されたバッチだけを処理します。

-- Trigger-based: only sees the rows inserted in the current batch
CREATE MATERIALIZED VIEW analytics.mv_realtime_counts
TO analytics.counts_table AS
SELECT pickup_date, count() AS trips
FROM default.trips_raw
GROUP BY pickup_date;

リフレッシュ型 MV(スケジュール実行): cron ジョブのように、スケジュールに従ってクエリ全体を再実行します。

-- Refreshable: runs the full SELECT every 3 minutes
CREATE MATERIALIZED VIEW analytics.mv_hourly_revenue
REFRESH EVERY 180 SECOND AS
SELECT
    toStartOfHour(pickup_at) AS hour_bucket,
    pickup_borough,
    sum(total_amount)        AS revenue
FROM analytics.fact_trips FINAL
GROUP BY hour_bucket, pickup_borough;

どちらを使うか:

  • トリガー型: 新しいデータだけを処理すればよい、挿入ストリーム上のリアルタイム集計
  • リフレッシュ型: fact_trips FINAL をクエリする集計(重複排除のためにテーブル全体を見る必要がある)、あるいはロジックの単純さと引き換えに数分の遅れを許容できるダッシュボード

リフレッシュ間隔の変更:

ALTER TABLE analytics.mv_hourly_revenue MODIFY REFRESH EVERY 60 SECOND;

2. ソートキー (ORDER BY)

Snowflake では、CLUSTER BY をオプティマイザへのヒントとして使いました。ClickHouse では ORDER BY が プライマリインデックス であり、ディスク上でのデータの物理的なソート順を決め、あらゆる範囲スキャンを左右します。

仕組み

ClickHouse は スパースなプライマリインデックス を保持します: 約8,192行(1つのデータグラニュール)ごとに1つのインデックスエントリです。ORDER BY の列で絞り込むと、ClickHouse はグラニュール全体を読まずに読み飛ばします。ClickHouse が毎秒何十億行もスキャンできるのはこのためです — データの大半はディスクから出ることさえありません。

カーディナリティの順序が重要

必ず カーディナリティの低い列を先に、高い列を最後に 置いてください。これにより、一般的なケースでインデックスの読み飛ばし能力が最大になります。

-- Good: low cardinality (borough, ~6 values) first, then high cardinality (trip_id)
ORDER BY (pickup_borough, toStartOfMonth(pickup_at), trip_id)

-- Bad: high cardinality first — the index can't skip anything useful
ORDER BY (trip_id, pickup_borough, pickup_at)

Snowflake の CLUSTER BY と ClickHouse の ORDER BY

項目Snowflake CLUSTER BYClickHouse ORDER BY
目的クエリ性能へのヒント物理的なソート順(必須)
適用バックグラウンドでの再クラスタリング(非同期)挿入時に常に適用
粒度マイクロパーティションデータグラニュール(約8,000行)
必須かいいえはい — すべての MergeTree テーブルに必要

例: 既存の Snowflake クラスタキーに合わせる

-- Snowflake
CLUSTER BY (DATE_TRUNC('month', PICKUP_AT), PICKUP_LOCATION_ID)

-- ClickHouse equivalent
ORDER BY (toStartOfMonth(pickup_at), pickup_location_id, trip_id)
-- Note: trip_id added as tiebreaker to ensure unique sort order

スキップインデックス(簡単に)

ORDER BY キーに含まれない列については、ClickHouse は スキップインデックス(bloom filter、minmax、set)をサポートしており、グラニュールごとに列レベルのメタデータを保持します。sort key の中でカーディナリティの高い列より後ろに来る、カーディナリティの低い列で絞り込む場合に有用です。

-- Add a bloom filter skip index on payment_type
ALTER TABLE analytics.fact_trips
ADD INDEX idx_payment_type payment_type TYPE bloom_filter GRANULARITY 4;

3. JSON の扱い

Snowflake の VARIANT 列型は、ネストした JSON をたどるためのコロンパス記法をサポートします。ClickHouse では代わりに、明示的な JSONExtract* 関数を使います。

対応表

SnowflakeClickHouse備考
col:key::FLOATJSONExtractFloat(col, 'key')トップレベルの float フィールド
col:driver.rating::FLOATJSONExtractFloat(col, 'driver', 'rating')ネストした float フィールド
col:app.surge_multiplier::FLOATJSONExtractFloat(col, 'app', 'surge_multiplier')ネストした float
col:route.waypoints[0]::STRINGJSONExtractString(col, 'route', 'waypoints', 0)インデックスで指定する配列要素
col:driver.id::INTJSONExtractInt(col, 'driver', 'id')整数フィールド

関数のバリエーション

-- Float (returns 0.0 if key missing or wrong type)
JSONExtractFloat(trip_metadata, 'driver', 'rating')

-- String (returns '' if missing)
JSONExtractString(trip_metadata, 'app', 'version')

-- Integer (returns 0 if missing)
JSONExtractInt(trip_metadata, 'driver', 'id')

-- Bool (returns 0/1)
JSONExtractBool(trip_metadata, 'app', 'is_shared')

-- Raw value as string (preserves JSON sub-object)
JSONExtractRaw(trip_metadata, 'route')

パフォーマンスのヒント

同じ JSON 列を繰り返しクエリするなら、下流のクエリごとに JSONExtractFloat を呼ぶのではなく、staging モデルの段階(stg_trips.sql の中)で型付きの列として取り出すことを検討してください。このラボの dbt モデルはそうしています。


4. 日付・時刻関数

Snowflake と ClickHouse の日付・時刻の機能は似ていますが、構文が異なります。最もよく使う変換は次のとおりです。

対応表

SnowflakeClickHouse備考
DATE_TRUNC('hour', col)toStartOfHour(col)時単位に切り捨て
DATE_TRUNC('day', col)toStartOfDay(col) または toDate(col)日単位に切り捨て
DATE_TRUNC('month', col)toStartOfMonth(col)月単位に切り捨て
CURRENT_TIMESTAMP()now()現在の日時
CURRENT_DATE()today()現在の日付
DATEADD('day', -7, CURRENT_DATE())today() - INTERVAL 7 DAY日付の演算
DATEDIFF('day', a, b)dateDiff('day', a, b)2つの日付の間の日数

ClickHouse にしかない便利な関数

yesterday()              -- today() - 1 day
toStartOfWeek(col)       -- Monday of the containing week
toStartOfQuarter(col)    -- first day of the quarter
toYear(col)              -- extract year as integer
toMonth(col)             -- extract month as integer (1-12)
toDayOfWeek(col)         -- 1=Monday, 7=Sunday

INTERVAL の構文

-- ClickHouse
now() - INTERVAL 7 DAY
now() - INTERVAL 1 HOUR
now() - INTERVAL 30 MINUTE
pickup_at + INTERVAL 90 SECOND

-- Snowflake equivalent
DATEADD('day', -7, CURRENT_TIMESTAMP())
DATEADD('hour', -1, CURRENT_TIMESTAMP())

5. 近似関数

ClickHouse は、何十億行に対する厳密な答えよりも、ダッシュボードに十分な精度の近似的な答えのほうが速い、という分析ワークロードのために作られています。ClickHouse には近似集計関数がいくつも組み込まれています。

個数のカウント(distinct)

関数精度速度使う場面
uniqExact(col)厳密最も遅いコンプライアンス報告、請求
uniq(col)誤差 約2%速いダッシュボード、探索的分析
uniqHLL12(col)誤差 約1.6%最速、メモリは 2.5KB 固定高カーディナリティ、メモリ制約がある場合
-- Exact (like Snowflake COUNT(DISTINCT ...))
SELECT uniqExact(trip_id) FROM analytics.fact_trips FINAL;

-- Approximate — good for "how many unique passengers today?"
SELECT uniq(passenger_id) FROM analytics.fact_trips FINAL;

パーセンタイル

関数備考
quantile(level)(col)厳密な分位点。メモリを多く使う
quantileTDigest(level)(col)t-digest による近似。メモリは固定
quantileTDigestWeighted(level)(col, weight)重み付きの t-digest
-- P95 trip duration — approximate but uses O(1) memory
SELECT quantileTDigest(0.95)(duration_minutes)
FROM analytics.fact_trips FINAL;

-- Multiple percentiles in one pass
SELECT quantileTDigestMerge(0.5)(state), quantileTDigestMerge(0.95)(state)
FROM analytics.fact_trips FINAL;

目安: インタラクティブなダッシュボードには uniq と quantileTDigest を使ってください。uniqExact と quantile は、課金、SLA、コンプライアンスで厳密な値が必要なときにだけ使ってください。


6. ディクショナリ

ディクショナリは インメモリのルックアップテーブル で、ClickHouse がホットな状態で保持し、クエリ時に事前結合済みの形で使えるようにします。フル JOIN のコストをかけずに結合したい小さなディメンションテーブルに相当する、ClickHouse 側の仕組みです。

何であるか

ディクショナリはソース(ClickHouse のテーブル、ファイル、または外部データベース)に支えられ、サービス起動時または SYSTEM RELOAD DICTIONARIES を呼んだときにメモリへロードされます。ルックアップはキー経由で行われ、1つ以上の属性を返します。

CREATE DICTIONARY の構文

-- From 04_create_dictionary.sql
CREATE DICTIONARY analytics.taxi_zones_dict (
    zone_id      UInt16,
    borough      String,
    service_zone String
)
PRIMARY KEY zone_id
SOURCE(CLICKHOUSE(
    TABLE 'dim_taxi_zones'
    DB    'analytics'
))
LIFETIME(MIN 300 MAX 600)   -- refresh every 5-10 minutes
LAYOUT(FLAT());             -- hash map, best for < 1M rows

LAYOUT の選択肢:

  • FLAT() — 整数キーで添字づけされた配列。最速だが、連続した整数キーが必要
  • HASHED() — ハッシュマップ。任意の整数キーで動作する
  • COMPLEX_KEY_HASHED() — 複合キーや文字列キーに対応するハッシュマップ

dictGet の使い方

-- Instead of: JOIN analytics.dim_taxi_zones USING (zone_id)
SELECT
    trip_id,
    dictGet('analytics.taxi_zones_dict', 'borough', toUInt64(pickup_location_id)) AS pickup_borough,
    dictGet('analytics.taxi_zones_dict', 'borough', toUInt64(dropoff_location_id)) AS dropoff_borough
FROM analytics.fact_trips FINAL;

ディクショナリと JOIN の使い分け

状況使うもの
小さく安定した参照テーブル(< 100万行、ほとんど変わらない)ディクショナリ
大きなディメンションテーブル、または頻繁に更新されるデータJOIN
同じルックアップを繰り返し実行するダッシュボードのクエリディクショナリ(初回ロード後はルックアップのコストがほぼゼロ)
一度きりの分析クエリJOIN

7. SAMPLE 句

ClickHouse はクエリ構文の中で行レベルのサンプリングを直接サポートします。サンプリングはデータの決定的な一部だけを読みます — 厳密な結果が不要な探索的分析に便利です。

構文

-- Read approximately 10% of rows
SELECT count(), avg(total_amount)
FROM analytics.fact_trips SAMPLE 0.1;

-- Read a specific number of rows (approximately)
SELECT trip_id, pickup_at, total_amount
FROM analytics.fact_trips SAMPLE 1000000;

結果のスケーリング

サンプリングするときは、集計値に 1 / sample_rate を掛けてテーブル全体の値を推定します。

-- Estimate total revenue from 10% sample
SELECT sum(total_amount) * 10 AS estimated_total_revenue
FROM analytics.fact_trips SAMPLE 0.1;

SAMPLE を使う場面

  • テーブル全体で実行する前の探索的分析(「クエリのロジックは合っているか」)
  • 近似値で構わないダッシュボードのタイル
  • 代表的なサブセットでの ML モデルの学習

注: SAMPLE は、テーブルの ORDER BY キーがサンプリング列で始まっていること、あるいは CREATE TABLE 文に SAMPLE BY 句を追加することを要求します。このラボの trips_raw テーブルは、この目的のために SAMPLE BY cityHash64(trip_id) を付けて作成されています。


8. dbt-clickhouse アダプターの注意点

dbt-clickhouse アダプター(dbt-clickhouse>=1.8)は標準的な dbt の機能をほぼサポートしますが、理解しておくべき ClickHouse 固有の挙動がいくつかあります。

delete_insert インクリメンタル戦略

ClickHouse には MERGE INTO がありません。dbt-clickhouse アダプターの delete_insert 戦略がそれをエミュレートします。

  1. 取り込むバッチとキー列が一致する行を、ターゲットテーブルから削除する
  2. 取り込む行をすべて挿入する
-- What dbt generates for incremental models
DELETE FROM analytics.fact_trips WHERE trip_id IN (SELECT trip_id FROM __dbt_tmp);
INSERT INTO analytics.fact_trips SELECT * FROM __dbt_tmp;

モデル側で次のように設定します。

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

engine と order_by の設定

すべての MergeTree テーブルにはエンジンと ORDER BY が必要です。dbt モデルの config で両方を指定してください。

{{
    config(
        engine='MergeTree()',
        order_by='(zone_id)'
    )
}}

cluster_by ではなく order_by

Snowflake の dbt モデルでは cluster_by を使っていたかもしれません。dbt-clickhouse では代わりに order_by を使います。ClickHouse に Snowflake の cluster_by に相当するものはありません — ORDER BY が常に物理的なソートです。

ClickHouse Cloud 向けの profiles.yml

ClickHouse Cloud は TLS を必須とします。secure: true を設定してください。

# ~/.dbt/profiles.yml
nyc_taxi_ch:
  target: dev
  outputs:
    dev:
      type: clickhouse
      host: "{{ env_var('CLICKHOUSE_HOST') }}"
      port: 8443
      user: default
      password: "{{ env_var('CLICKHOUSE_PASSWORD') }}"
      schema: analytics       # default database for models without a custom schema
      secure: true
      threads: 4

スキーマの命名と別データベース

Snowflake が schema と呼ぶものを、ClickHouse では database と呼びます。dbt-clickhouse アダプターは dbt のスキーマを ClickHouse のデータベースにマッピングします。このプロジェクトの generate_schema_name マクロは dbt のデフォルトの挙動を上書きし、+schema: analytics を持つモデルが staging_analytics ではなく analytics データベースに配置されるようにしています。

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

これはパート1の Snowflake dbt プロジェクトで使ったのと同じパターンです — 両方のアダプターでスキーマ命名の挙動を揃えるため、マクロは意図的に同一にしてあります。

このページの内容

JA