このチュートリアルでは、materialized views を使用して、大量のイベントテーブルから事前集計済みのロールアップを維持する方法を紹介します。 作成するオブジェクトは 3 つです。生テーブル、ロールアップテーブル、そしてロールアップへ自動的に書き込むmaterialized viewです。
このパターンを使う場面
次のような場合は、このパターンを使用します。
- 追記専用のイベントストリーム (クリック、ページビュー、IoT、ログ) がある。
- クエリの大半が、一定の時間範囲に対する集計 (分単位/時間単位/日単位) である。
- 生のすべての行を毎回再スキャンすることなく、安定してサブ秒で読み取りたい。
raw eventsテーブルを作成する
CREATE TABLE events_raw
(
event_time DateTime,
user_id UInt64,
country LowCardinality(String),
event_type LowCardinality(String),
value Float64
)
ENGINE = MergeTree
PARTITION BY toYYYYMM(event_time)
ORDER BY (event_time, user_id)
TTL event_time + INTERVAL 90 DAY DELETE注記
PARTITION BY toYYYYMM(event_time)を使うと、パーティションを小さく保てるため、削除しやすくなります。ORDER BY (event_time, user_id)は、時間範囲を絞ったクエリに加え、二次フィルタにも対応できます。LowCardinality(String)は、カテゴリ次元のメモリ使用量を削減します。TTLにより、90日後に生データがクリーンアップされます (保持要件に合わせて調整してください) 。
ロールアップ (集計済み) テーブルを設計する
時間単位の粒度で事前集計します。 粒度は、最も一般的な分析期間に合わせて選択してください。
CREATE TABLE events_rollup_1h
(
bucket_start DateTime, -- start of the hour
country LowCardinality(String),
event_type LowCardinality(String),
users_uniq AggregateFunction(uniqExact, UInt64),
value_sum AggregateFunction(sum, Float64),
value_avg AggregateFunction(avg, Float64),
events_count AggregateFunction(count)
)
ENGINE = AggregatingMergeTree
PARTITION BY toYYYYMM(bucket_start)
ORDER BY (bucket_start, country, event_type)aggregate states (例: AggregateFunction(sum, ...)) を保存します。これは部分的な集計結果をコンパクトに表現するもので、後からマージしたり最終化したりできます。
ロールアップを生成する materialized view を作成する
この materialized view は events_raw への挿入時に自動的にトリガーされ、集計状態をロールアップに書き込みます。
CREATE MATERIALIZED VIEW mv_events_rollup_1h
TO events_rollup_1h
AS
SELECT
toStartOfHour(event_time) AS bucket_start,
country,
event_type,
uniqExactState(user_id) AS users_uniq,
sumState(value) AS value_sum,
avgState(value) AS value_avg,
countState() AS events_count
FROM events_raw
GROUP BY bucket_start, country, event_type;サンプルデータを挿入する
サンプルデータを挿入します。
INSERT INTO events_raw VALUES
(now() - INTERVAL 4 SECOND, 101, 'US', 'view', 1),
(now() - INTERVAL 3 SECOND, 101, 'US', 'click', 1),
(now() - INTERVAL 2 SECOND, 202, 'DE', 'view', 1),
(now() - INTERVAL 1 SECOND, 101, 'US', 'view', 1);ロールアップのクエリ
集約状態は、読み取り時にマージすることも、確定することもできます。
SELECT
bucket_start,
country,
event_type,
uniqExactMerge(users_uniq) AS users,
sumMerge(value_sum) AS value_sum,
avgMerge(value_avg) AS value_avg,
countMerge(events_count) AS events
FROM events_rollup_1h
WHERE bucket_start >= now() - INTERVAL 1 DAY
GROUP BY ALL
ORDER BY bucket_start, country, event_type;SELECT
bucket_start,
country,
event_type,
uniqExactMerge(users_uniq) AS users,
sumMerge(value_sum) AS value_sum,
avgMerge(value_avg) AS value_avg,
countMerge(events_count) AS events
FROM events_rollup_1h
WHERE bucket_start >= now() - INTERVAL 1 DAY
GROUP BY ALL
ORDER BY bucket_start, country, event_type
SETTINGS final = 1; -- または SELECT ... FINAL を使用最適なパフォーマンスを得るには、主キー内のフィールドでフィルタリングします
EXPLAIN コマンドを使うと、索引を使ってどのようにデータが絞り込まれるかを確認できます。
EXPLAIN indexes=1
SELECT *
FROM events_rollup_1h
WHERE bucket_start BETWEEN now() - INTERVAL 3 DAY AND now()
AND country = 'US'; ┌─explain────────────────────────────────────────────────────────────────────────────────────────────────────────────────────────────┐
1. │ Expression ((Project names + Projection)) │
2. │ Expression │
3. │ ReadFromMergeTree (default.events_rollup_1h) │
4. │ Indexes: │
5. │ MinMax │
6. │ Keys: │
7. │ bucket_start │
8. │ Condition: and((bucket_start in (-Inf, 1758550242]), (bucket_start in [1758291042, +Inf))) │
9. │ Parts: 1/1 │
10. │ Granules: 1/1 │
11. │ Partition │
12. │ Keys: │
13. │ toYYYYMM(bucket_start) │
14. │ Condition: and((toYYYYMM(bucket_start) in (-Inf, 202509]), (toYYYYMM(bucket_start) in [202509, +Inf))) │
15. │ Parts: 1/1 │
16. │ Granules: 1/1 │
17. │ PrimaryKey │
18. │ Keys: │
19. │ bucket_start │
20. │ country │
21. │ Condition: and((country in ['US', 'US']), and((bucket_start in (-Inf, 1758550242]), (bucket_start in [1758291042, +Inf)))) │
22. │ Parts: 1/1 │
23. │ Granules: 1/1 │
└────────────────────────────────────────────────────────────────────────────────────────────────────────────────────────────────────┘上記のクエリ実行計画では、3種類の索引が使われていることがわかります:
MinMax 索引、パーティション索引、主キー索引です。
各索引は、主キーで指定したフィールド (bucket_start, country, event_type) を利用しています。
フィルタリング性能を最大限に引き出すには、クエリで主キーのフィールドを使ってデータを絞り込めるようにしてください。
よくあるバリエーション
- 異なる粒度: 日次のロールアップを追加します:
CREATE TABLE events_rollup_1d
(
bucket_start Date,
country LowCardinality(String),
event_type LowCardinality(String),
users_uniq AggregateFunction(uniqExact, UInt64),
value_sum AggregateFunction(sum, Float64),
value_avg AggregateFunction(avg, Float64),
events_count AggregateFunction(count)
)
ENGINE = AggregatingMergeTree
PARTITION BY toYYYYMM(bucket_start)
ORDER BY (bucket_start, country, event_type);次に、2つ目のmaterialized view:
CREATE MATERIALIZED VIEW mv_events_rollup_1d
TO events_rollup_1d
AS
SELECT
toDate(event_time) AS bucket_start,
country,
event_type,
uniqExactState(user_id),
sumState(value),
avgState(value),
countState()
FROM events_raw
GROUP BY ALL;- 圧縮: rawテーブル上の大きなカラムにコーデック (例:
Codec(ZSTD(3))) を適用します。 - コスト管理: 保持コストの高いデータ保持はrawテーブル側に寄せ、長期間保持するロールアップは維持します。
- バックフィル: 履歴データを読み込む際は、
events_rawに insert して、materialized view でロールアップを自動的に構築させます。既存の行に対しては、適切であれば materialized view 作成時にPOPULATEを使うか、INSERT SELECTを使用します。
クリーンアップと保持期間
- 生データの有効期限 (TTL) は延ばし (例: 30/90 日) 、ロールアップはより長く保持します (例: 1 年) 。
- 階層化が有効になっている場合は、TTL to move を使って古いパーツをより低コストなストレージに移動することもできます。
トラブルシューティング
- materialized view が更新されない場合は、挿入先が events_raw (ロールアップテーブルではなく) になっていること、および materialized view のターゲットが正しいこと (
TO events_rollup_1h) を確認してください。 - クエリが遅い場合は、ロールアップが使われていること (ロールアップテーブルを直接クエリする) と、時間フィルターがロールアップの粒度に合っていることを確認してください。
- バックフィルに不整合がある場合は、
SYSTEM FLUSH LOGSを実行し、system.query_log/system.partsを確認して、insert と merge が行われていることを確かめてください。