•19 min read

ClickHouseのMaterializedViewとReplacingMergeTree:サブ秒リアルタイム分析

ClickHouseのMaterializedViewとReplacingMergeTree:サブ秒リアルタイム分析

ClickHouseはリアルタイム分析ワークロードに優れています。カーディナリティが高く、大量のデータストリームに対してサブ秒のクエリレイテンシーを達成するには、多くの場合、事前集計と効率的な重複排除が必要です。このガイドでは、堅牢なリアルタイム分析パイプラインを構築するために、AggregatingMergeTreeおよびReplacingMergeTreeエンジンを使用したClickHouseマテリアライズドビューのアーキテクチャと実装について詳しく説明します。

Audio Briefing
0:00 / 0:00

アーキテクチャの概要:ストリーム集計と重複排除

ここで取り組む主要な問題は、イベントが順不同で到着したり重複したりする可能性があるイベントストリームから、リアルタイムで集計されたメトリクスが必要であるということです。私たちのソリューションには以下が含まれます。

  1. 生イベントテーブル: すべての受信イベントを保存する追記専用テーブル。これは信頼できる情報源として機能します。
  2. 重複排除レイヤー: 遅れて到着したイベントや再実行されたイベントを処理し、イベントの冪等性を保証するReplacingMergeTreeテーブル。
  3. 集計用マテリアライズドビュー: マテリアライズドビューによってデータが投入されるAggregatingMergeTreeテーブルで、集計を事前に計算します。このテーブルはステートフルな集計を保存し、クエリ時間を大幅に短縮します。

この階層化されたアプローチは、データの整合性とクエリパフォーマンスの両方を提供します。

Advertisement

ReplacingMergeTreeによる生イベントの取り込み

まず、生イベントテーブルを定義します。単純なMergeTreeでも十分ですが、特にイベントの再実行やKafkaのようなメッセージキューからの少なくとも1回の配信セマンティクスを扱う場合、冪等な取り込みにはReplacingMergeTreeが不可欠です。

CREATE TABLE IF NOT EXISTS events_raw (
    event_id UUID,
    user_id String,
    event_type LowCardinality(String),
    event_time DateTime64(3),
    value Float64,
    _version UInt64 DEFAULT 1 -- Version column for ReplacingMergeTree
) ENGINE = ReplacingMergeTree(_version)
ORDER BY (event_id, event_time)
PRIMARY KEY (event_id);
  • ReplacingMergeTree(_version): このエンジンは、特定のORDER BYキー(ここではevent_id, event_time)に対して、マージ中に最大_versionを持つ行のみが保持されることを保証します。_versionが省略された場合、挿入順序で最後の行が保持されます。
  • ORDER BY (event_id, event_time): ソートキーを定義します。ReplacingMergeTreeはこれを使用して「重複」行を識別します。
  • PRIMARY KEY (event_id): event_idに対するポイントルックアップと範囲スキャンを最適化します。

重複するevent_idとより高いバージョンを含むサンプルデータをいくつか挿入してみましょう。

INSERT INTO events_raw (event_id, user_id, event_type, event_time, value, _version) VALUES
('a0000000-0000-0000-0000-000000000001', 'user1', 'page_view', '2023-10-26 10:00:00.000', 1.0, 1),
('a0000000-0000-0000-0000-000000000002', 'user1', 'click', '2023-10-26 10:00:05.000', 0.5, 1),
('a0000000-0000-0000-0000-000000000003', 'user2', 'page_view', '2023-10-26 10:00:10.000', 1.0, 1),
('a0000000-0000-0000-0000-000000000001', 'user1', 'page_view', '2023-10-26 10:00:00.000', 1.2, 2); -- Duplicate event_id, higher version

重複排除を観察するには、マージを強制するか、バックグラウンドマージを待つ必要があります。

OPTIMIZE TABLE events_raw FINAL;

SELECT event_id, user_id, value, _version FROM events_raw ORDER BY event_id;

出力:

┌─event_id─────────────────────────────┬─user_id─┬─value─┬─_version─┐
│ a0000000-0000-0000-0000-000000000001 │ user1   │   1.2 │        2 │
│ a0000000-0000-0000-0000-000000000002 │ user1   │   0.5 │        1 │
│ a0000000-0000-0000-0000-000000000003 │ user2   │   1.0 │        1 │
└──────────────────────────────────────┴─────────┴───────┴──────────┘

event_id a0000000-0000-0000-0000-000000000001がvalue 1.2と_version 2になったことに注目してください。これは重複排除を示しています。

リアルタイム集計のためのマテリアライズドビュー

ClickHouseのマテリアライズドビューは、従来のRDBMSのような事前計算されたスナップショットではありません。代わりに、ソーステーブルに挿入された新しいデータに対してSELECTクエリを実行し、その結果をターゲットテーブルに書き込むトリガーです。これにより、継続的かつ増分的な集計に最適です。

ステートフル集計のためのAggregatingMergeTree

マテリアライズドビューのターゲットテーブルはAggregatingMergeTreeエンジンを使用します。このエンジンは、集計関数の状態を保存し、最終的な値は保存しません。データパーツがマージされると、これらの状態は*Merge関数を使用して結合されます。

集計テーブルを定義します。

CREATE TABLE IF NOT EXISTS events_agg_mv (
    event_date Date,
    event_type LowCardinality(String),
    user_id String,
    total_value AggregateFunction(sum, Float64),
    unique_users AggregateFunction(uniqHLL12, String),
    event_count AggregateFunction(count)
) ENGINE = AggregatingMergeTree()
ORDER BY (event_date, event_type, user_id);
  • AggregateFunction(sum, Float64): この特殊なデータ型は、sum集計関数の中間状態を保存します。
  • uniqHLL12: カーディナリティの高いデータに適した、非常に効率的な近似個別カウントアルゴリズムです。
  • ORDER BY (event_date, event_type, user_id): 集計キーを定義します。同じキーを持つすべての行がマージされ、その集計状態が結合されます。

マテリアライズドビューの作成

次に、events_rawからevents_agg_mvにデータを投入するマテリアライズドビューを作成します。

CREATE MATERIALIZED VIEW IF NOT EXISTS mv_events_agg
TO events_agg_mv
AS SELECT
    toDate(event_time) AS event_date,
    event_type,
    user_id,
    sumState(value) AS total_value,
    uniqHLL12State(user_id) AS unique_users,
    countState() AS event_count
FROM events_raw
GROUP BY event_date, event_type, user_id;
  • TO events_agg_mv: ターゲットテーブルを指定します。
  • sumState(value), uniqHLL12State(user_id), countState(): これらは集計関数の「状態」バージョンです。集計の中間状態を返し、AggregatingMergeTreeがそれを保存します。

events_rawにさらにデータを挿入し、マテリアライズドビューの効果を観察してみましょう。

INSERT INTO events_raw (event_id, user_id, event_type, event_time, value, _version) VALUES
('a0000000-0000-0000-0000-000000000004', 'user1', 'page_view', '2023-10-26 10:00:15.000', 1.0, 1),
('a0000000-0000-0000-0000-000000000005', 'user2', 'click', '2023-10-26 10:00:20.000', 0.8, 1),
('a0000000-0000-0000-0000-000000000006', 'user1', 'page_view', '2023-10-27 11:00:00.000', 1.0, 1);

集計テーブルをクエリします。集計を確定するために*Merge関数を使用していることに注目してください。

SELECT
    event_date,
    event_type,
    user_id,
    sumMerge(total_value) AS final_total_value,
    uniqHLL12Merge(unique_users) AS final_unique_users,
    countMerge(event_count) AS final_event_count
FROM events_agg_mv
GROUP BY event_date, event_type, user_id
ORDER BY event_date, event_type, user_id;

出力:

┌─event_date─┬─event_type─┬─user_id─┬─final_total_value─┬─final_unique_users─┬─final_event_count─┐
│ 2023-10-26 │ click      │ user1   │               0.5 │                  1 │                 1 │
│ 2023-10-26 │ click      │ user2   │               0.8 │                  1 │                 1 │
│ 2023-10-26 │ page_view  │ user1   │               2.2 │                  1 │                 2 │
│ 2023-10-26 │ page_view  │ user2   │               1.0 │                  1 │                 1 │
│ 2023-10-27 │ page_view  │ user1   │               1.0 │                  1 │                 1 │
└────────────┴────────────┴─────────┴───────────────────┴────────────────────┴───────────────────┘

集計はリアルタイムで正しく更新されます。ReplacingMergeTree on events_rawは、イベントがより高いバージョンで再送信された場合、マテリアライズドビューが更新されたイベントを処理し、その後のマージでAggregatingMergeTreeが変更を正しく反映することを保証します。

比較:通常ビュー vs. マテリアライズドビュー

機能通常ビュー(例:CREATE VIEW)マテリアライズドビュー(例:CREATE MATERIALIZED VIEW)
データストレージデータは保存されず、クエリはオンデマンドで実行されるデータは事前に計算され、ターゲットテーブルに保存される
クエリパフォーマンス基になるテーブルに依存し、複雑な集計では遅くなる可能性がある事前集計されたデータでは非常に高速、サブ秒
データ鮮度常にリアルタイムリアルタイム(データがソーステーブルに挿入されると)
リソース使用量ストレージは低いが、クエリCPU/IOは高いストレージは高いが、クエリCPU/IOは低い
ユースケース単純なエイリアス、複雑なアドホッククエリリアルタイムダッシュボード、固定レポート、大量分析

50ミリ秒未満の分析クエリの最適化

50ミリ秒未満のレイテンシーを達成するには、テーブル設計、クエリパターン、ClickHouse構成を慎重に検討する必要があります。

  1. AggregatingMergeTree ORDER BYキー: events_agg_mvのORDER BY句は非常に重要です。分析クエリの一般的なGROUP BYおよびWHERE句と一致させる必要があります。たとえば、event_dateとevent_typeで頻繁にクエリを実行する場合、これらはORDER BYの先頭列である必要があります。
  2. PRIMARY KEY: AggregatingMergeTreeの場合、PRIMARY KEYは通常、ORDER BYキーのプレフィックスです。これにより、データパーツを迅速にプルーニングできます。
  3. LowCardinalityデータ型: 識別可能な値の数が限られている列(例:event_type)にはLowCardinality(String)を使用します。これにより、ストレージが大幅に削減され、辞書エンコーディングによりクエリパフォーマンスが向上します。
  4. DateTime64 vs DateTime: DateTime64はミリ秒単位の精度を提供し、イベントストリームにはしばしば必要です。toDate()またはtoStartOfHour()関数が目的の集計粒度と一致していることを確認してください。
  5. FINALキーワード: AggregatingMergeTreeまたはReplacingMergeTreeをクエリする場合、SELECT ... FROM table FINALはすべてのマージが完了し、完全に集計/重複排除された結果が得られることを保証します。ただし、FINALはマージを強制するため遅くなる可能性があります。リアルタイムダッシュボードの場合、わずかに古いデータを許容し、FINALを省略してバックグラウンドマージに頼ることもできます。重要なレポートの場合、FINALが必要です。
  6. index_granularity: この設定(デフォルト8192)は、インデックス作成用のデータブロック内の行数を決定します。調整することでパフォーマンスに影響を与える可能性がありますが、デフォルトで十分な場合が多いです。
  7. ハードウェア: ClickHouseのパフォーマンスには、十分なRAM、高速NVMe SSD、CPUコアが不可欠です。
  8. 分散テーブル: 非常に大規模なデータセットの場合、AggregatingMergeTreeの上にDistributedテーブルを使用して水平にスケーリングします。
Advertisement

本番環境での落とし穴とトラブルシューティング

  1. マテリアライズドビューの遅延:

    • 症状: events_agg_mvの更新が遅い、またはそれに対するクエリが古いデータを表示する。
    • 原因: events_rawへの取り込みレートが高いことと、複雑なMVロジックまたはリソース制約の組み合わせ。マテリアライズドビューは挿入時にデータを同期的に処理します。MVクエリが遅い場合、挿入がブロックされる可能性があります。
    • 修正:
      • MV SELECTクエリを簡素化します。
      • 効率的なMV処理のために、events_rawに適切なORDER BYとPRIMARY KEYがあることを確認します。
      • ClickHouseのリソース(CPU、RAM、IO)をスケールアップします。
      • 非同期マテリアライズドビューの使用を検討します(TO target_tableを省略し、MVに独自の.テーブルを作成させ、その後別のプロセスまたは別のMVを介して.テーブルからAggregatingMergeTreeテーブルに挿入します)。これにより、取り込みと集計が分離されますが、複雑さが増します。ほとんどの場合、同期MVはシンプルさとリアルタイム保証のために推奨されます。
  2. ReplacingMergeTreeが重複排除しない:

    • 症状: OPTIMIZE TABLE FINAL後も重複するevent_idが残る。
    • 原因: ORDER BYキーが正しくない。ReplacingMergeTreeはORDER BYキーに基づいて重複排除を行います。event_idがORDER BYキーの一部でない場合、または_version列が正しく使用されていない場合、期待どおりに重複排除は行われません。
    • 修正: ORDER BYに一意の識別子(例:event_id)とバージョン列(例:_version)が含まれていること、およびエンジン定義で正しく指定されていることを確認します。
  3. AggregatingMergeTreeクエリのパフォーマンス:

    • 症状: events_agg_mvに対するクエリが、*Merge関数を使用しても遅い。
    • 原因:
      • クエリのGROUP BY句がevents_agg_mvのORDER BYキーと一致しない。これにより、ClickHouseは必要以上に多くのデータを読み込むことになります。
      • GROUP BY列の個別値が多すぎるため、多数の小さなデータパーツが発生したり、マージ中にメモリ使用量が増加したりする。
      • FINALを不必要に使用している。
    • 修正:
      • 一般的なクエリパターンに合わせてevents_agg_mv ORDER BYをリファクタリングします。
      • PRIMARY KEYがORDER BYのプレフィックスであることを確認します。
      • 正確性のために厳密に必要でない限り、FINALの使用を避けます。
      • 中間集計がまだ細かすぎる場合は、さらに事前集計を検討します。
  4. ディスクスペースの消費:

    • 症状: AggregatingMergeTreeテーブルが過剰なディスクスペースを消費する。
    • 原因: AggregateFunctionの状態は、特にuniqHLL12StateやgroupArrayStateのような関数では、最終的な値よりも大きくなる可能性があります。また、ORDER BYキーのカーディナリティが高い場合、多くの小さなデータパーツにつながる可能性があります。
    • 修正:
      • 集計関数を見直します。uniqHLL12は個別カウントに対してスペース効率が良いです。groupArrayStateは非常に大きくなる可能性があります。
      • 集計粒度とデータパーツサイズのバランスを取るようにORDER BYキーを選択します。
      • events_rawとevents_agg_mvにTTL(Time-To-Live)ポリシーを実装し、古いデータを自動的に削除します。
-- Example TTL for events_raw (delete raw data after 30 days)
ALTER TABLE events_raw MODIFY TTL event_time + INTERVAL 30 DAY;

-- Example TTL for events_agg_mv (delete aggregates after 365 days)
ALTER TABLE events_agg_mv MODIFY TTL event_date + INTERVAL 365 DAY;

よくある質問

  1. マテリアライズドビューは作成後に変更できますか? いいえ、マテリアライズドビューは直接変更できません。DROPして再度CREATEする必要があります。そのため、MVを慎重に設計することが重要です。ターゲットテーブルのスキーマが変更された場合も、MVを削除して再作成する必要があります。

  2. マテリアライズドビューのソーステーブルからデータが削除された場合、どうなりますか? マテリアライズドビューはINSERT操作にのみ反応します。ソーステーブルでのDELETEまたはUPDATE操作は、マテリアライズドビューのターゲットテーブルに自動的に伝播しません。削除を処理する必要がある場合は、通常、「ソフト削除」メカニズム(例:is_deletedフラグ)を実装してクエリでフィルタリングするか、より複雑なCollapsingMergeTreeまたはVersionedCollapsingMergeTree設定を使用します。

  3. ReplacingMergeTreeは同じORDER BYキーを持つ同時挿入をどのように処理しますか? ReplacingMergeTreeは、マージ中に置換ロジックを適用することで同時挿入を処理します。同じORDER BYキーだが異なる_version値を持つ2つの挿入が同時に到着した場合、最初は別々のデータパーツに存在します。これらのパーツがマージされると、最も高い_versionを持つ行が保持されます。これにより、最終的な整合性が保証されます。

  4. AggregatingMergeTreeと通常のMergeTreeにGROUP BYを使用する場合の使い分けは? リアルタイムでデータを継続的に集計し、その集計を頻繁にクエリする必要がある場合は、AggregatingMergeTreeを使用します。これは集計状態を事前に計算して保存するため、クエリが大幅に高速になります。リアルタイムパフォーマンスが重要でない場合、または集計キーが非常に動的で事前に定義できない場合のアドホックな生データ集計には、MergeTreeとGROUP BYを使用します。AggregatingMergeTreeは、固定された既知の集計パターン向けです。

  5. マテリアライズドビューはすべての種類の集計に適していますか? マテリアライズドビューは、加算可能な集計、またはAggregateFunction状態を使用して表現できる集計に最適です。これには、sum、count、min、max、uniqHLL12、avg(sumStateとcountStateを使用)などが含まれます。すべての生行へのアクセスを必要とする集計(例:特定のAggregateFunctionサポートなしのquantile、median)や複雑なウィンドウ関数は、通常、直接的なマテリアライズドビューの事前集計には適しておらず、生データまたはより粒度の高い集計で実行する方が良いでしょう。

Share this article:

Stay Updated

Get the latest posts delivered straight to your inbox.

Free Developer Utilities

Free In-Browser Developer Tools

Clean AI CLI logs, build cron expressions, decode JWTs, and calculate chmod permissions offline.

Explore Tools
Advertisement