要約表を分けて挿入し、合計と状態で元データと同じ答えを出す
目標
元のイベントから(site, day)の集計をSummingMergeTreeとAggregatingMergeTreeに分けて入れ、マージ前でも元データとまったく同じ答えを出すクエリを書きます。合計で減らせない値は状態として持つ必要があること、その状態をより大きな単位でもう一度まとめる方法を確認します。
なぜ重要なのか
集計テーブルはダッシュボードを安くしますが、マージが終わる前は、同じキーが複数の行に散らばっています。その状態でSELECT *で読んだり、異なるユーザー数を数値で保存して足したりすると、間違った答えが黙って出ます。このラボではマージを止めてその間違った状態を固定し、採点ツールが、保存されたクエリの結果を、元データagg.hitsを直接集計した値と列ごとに突き合わせます。何回に分けて入れたかは、system.part_logの記録で確認します。
ステップ
- データベース
aggとテーブルagg.hitsを作成してください。列はts DateTime, site LowCardinality(String), user_id UInt64, dur_ms UInt32, is_bot UInt8(この順序)、エンジンはMergeTree、ORDER BY (site, ts)です。 /opt/lab/fixtures/aggregating/hits.sqlを1回だけ実行して、100万行を入れてください。agg.daily (site LowCardinality(String), day Date, hits UInt64, dur_total UInt64)をENGINE = SummingMergeTree ORDER BY (site, day)で作成し、SYSTEM STOP MERGES agg.dailyでマージを止めたあと、agg.hitsを時間帯ごと(toHour(ts) < 8、8 이상 16 미만、16 이상。2つ目のコードは韓国語で「8以上16未満」を、3つ目のコードは韓国語で「16以上」を意味する表記です)に3回のINSERTに分けて、site, toDate(ts), count(), sum(dur_ms)を入れてください。- マージ前の今でも、元データと同じ(site, day)ごとのビュー数・滞在時間の合計を出すクエリ(/root/ch/aggregating/q_daily.sql)を作成してください(agg.dailyだけを読み、列は
site, day, hits, dur_totalの順序)。そして、agg.dailyの現在の行数と、異なる(site, day)の数を、raw_rows・keysとして書き込んでください(保存先: /root/ch/aggregating/raw.json)。 SYSTEM START MERGES agg.dailyのあと、OPTIMIZE TABLE agg.daily FINALでまとめて、キーごとに1行になるようにしてください。agg.daily_state (site LowCardinality(String), day Date, users AggregateFunction(uniqExact, UInt64), avg_dur AggregateFunction(avg, UInt32), human_hits AggregateFunction(countIf, UInt8))をENGINE = AggregatingMergeTree ORDER BY (site, day)で作成してマージを止めたあと、ステップ3と同じ時間帯の3つのまとまりでuniqExactState(user_id)・avgState(dur_ms)・countIfState(is_bot = 0)を入れてください。- agg.daily_stateだけを読んで、(site, day)ごとの
users・avg_dur・human_hitsを出すクエリ(/root/ch/aggregating/q_state.sql)を作成してください(列はsite, day, users, avg_dur, human_hitsの順序)。そして、マージ前の各行のfinalizeAggregation(users)をそのまま足した値と、q_state.sqlのusersをすべて足した値を、naive_users_sum・true_users_sumとして書き込んでください(保存先: /root/ch/aggregating/naive.json)。 agg.day_total (day Date, users AggregateFunction(uniqExact, UInt64))をENGINE = AggregatingMergeTree ORDER BY dayで作成し、agg.daily_stateの状態をuniqExactMergeState(users)で日付ごとにまとめ直して入れてください。そして、agg.day_totalだけを読んで、日付ごとの(サイトの区別がない)異なるユーザー数を出すクエリ(/root/ch/aggregating/q_day.sql)を作成してください(列はday, users)。
参考
- サーバーはPodの起動時にすでに立ち上がっています。
clickhouse-clientと入力するだけで接続できます。止まっていた場合はch-upを実行してください(サーバーを立ち上げ直すとSTOP MERGESが解除されます)。 - INSERTの記録は
system.part_log(event_type = 'NewPart')に残ります。採点ツールは、テーブルのUUIDで絞り込みます。やり直すには、TRUNCATE TABLEしてから3回入れます。 - 状態の列は、人間が読める値ではありません。確認するときは、
finalizeAggregation(열)(プレースホルダーは列名です)や-Merge関数で仕上げて見ます。 - よくある間違い: 1回のINSERTで元データ全体を入れることです。同じキーが入れた瞬間にまとめられて、「マージ前」が見えなくなります。もう1つは、クエリで
sum(uniqExactMerge(...))のように、仕上げた値をもう一度足すことです。重なった人を何度も数えます。 - 公式ドキュメント: SummingMergeTree・AggregatingMergeTree・Aggregate Function Combinators・system.part_log
元のイベントテーブル
データベースaggとテーブルagg.hitsを作成してください。列はts DateTime, site LowCardinality(String), user_id UInt64, dur_ms UInt32, is_bot UInt8の順序、エンジンはMergeTree、ソートキーはORDER BY (site, ts)です。
集計テーブルは、いつも元データから作られます。元データをMergeTreeのまま置いておけば、集計を間違えて作っても、作り直せます。ドキュメントがSummingMergeTreeをMergeTreeと一緒に使うよう勧めている理由です。
元データ100万行
/opt/lab/fixtures/aggregating/hits.sqlを1回だけ実行して、agg.hitsに1,000,000行を入れてください。
clickhouse-clientに--queries-fileで渡せば実行できます。サイト5個×2026年9月の30日分なので、(site, day)の組み合わせは150個です。2回入れてしまった場合は、TRUNCATEしてからもう一度入れてください。
SummingMergeTreeに3回に分けて入れる
agg.daily (site LowCardinality(String), day Date, hits UInt64, dur_total UInt64)をENGINE = SummingMergeTree ORDER BY (site, day)で作成し、SYSTEM STOP MERGES agg.dailyでマージを止めたあと、agg.hitsを時間帯ごと(toHour(ts) < 8・8 이상 16 미만・16 이상。2つ目のコードは韓国語で「8以上16未満」を、3つ目のコードは韓国語で「16以上」を意味する表記です)に3回のINSERTに分けて、site, toDate(ts), count(), sum(dur_ms)を(site, day)でグループ化して入れてください。
1日分の集計が、3回に分かれて届く状況です。各INSERTはその時間帯だけをGROUP BYして入れるので、マージ前には、同じ(site, day)が3行あります。マージを止めるのは、入れる前に行う必要があります。
マージ前でも正しいクエリ
agg.dailyだけを読んで、(site, day)ごとのビュー数・滞在時間の合計を、元データとまったく同じに出すクエリ(/root/ch/aggregating/q_daily.sql)を作成してください。列はsite, day, hits, dur_totalの順序です。そして、agg.dailyの現在の行数と、異なる(site, day)の数を、raw_rows・keysとして書き込んでください(保存先: /root/ch/aggregating/raw.json)。
マージが終わったかどうかはわからないので、クエリの中でもう一度グループ化して足します。SELECT *で読むと、今はキーごとに3行が出ます。元データのagg.hitsを読んではいけません。集計テーブルだけで答えを出すのが目的です。
マージするとキーごとに1行になる
SYSTEM START MERGES agg.dailyのあと、OPTIMIZE TABLE agg.daily FINALでまとめて、agg.dailyがパート1つ・キーごとに1行になるようにしてください。
マージが、同じキーの数値列を足して1行に畳みます。まとめたあと、SELECT *が元データの集計と同じになるかを確認してください。本番環境ではマージのタイミングがわからないので、ステップ4のようなクエリが、やはり必要です。
合計で減らせない値は状態にする
agg.daily_state (site LowCardinality(String), day Date, users AggregateFunction(uniqExact, UInt64), avg_dur AggregateFunction(avg, UInt32), human_hits AggregateFunction(countIf, UInt8))をENGINE = AggregatingMergeTree ORDER BY (site, day)で作成し、SYSTEM STOP MERGES agg.daily_stateのあと、ステップ3と同じ時間帯の3つのまとまりでuniqExactState(user_id)・avgState(dur_ms)・countIfState(is_bot = 0)を(site, day)でグループ化して入れてください。
列の型AggregateFunction(関数, 引数の型)は、その関数の中間状態を持つという意味で、入れるときは、同じ関数に-Stateを付けて作ります。-Ifコンビネーターは、条件を最後の引数として受け取ります。マージ前なので、3つのまとまりの状態が別々にある必要があります。
状態はまとめてから仕上げる
agg.daily_stateだけを読んで、(site, day)ごとのusers(異なるユーザー数)・avg_dur(平均滞在時間)・human_hits(ボットではないビュー数)を、元データとまったく同じに出すクエリ(/root/ch/aggregating/q_state.sql)を作成してください(列の順序はsite, day, users, avg_dur, human_hits)。そして、マージ前の各行のfinalizeAggregation(users)をそのまま足した値と、q_state.sqlのusersをすべて足した値を、naive_users_sum・true_users_sumとして書き込んでください(保存先: /root/ch/aggregating/naive.json)。
-Mergeを付けた関数が、同じキーの状態をまとめてから結果を出します。finalizeAggregationは、状態1つをその場で仕上げるので、行ごとに仕上げた数値を足すと、複数の時間帯に来た人を何度も数えてしまいます。2つの合計がどれだけ違うかを見てください。
状態をより大きな単位でもう一度まとめる
agg.day_total (day Date, users AggregateFunction(uniqExact, UInt64))をENGINE = AggregatingMergeTree ORDER BY dayで作成し、agg.daily_stateのusersの状態をuniqExactMergeState(users)で日付ごとにまとめて入れ、agg.day_totalだけを読んで、日付ごと(サイトの区別なし)の異なるユーザー数を出すクエリ(/root/ch/aggregating/q_day.sql)を作成してください(列はday, users)。
-MergeStateは、状態をまとめて、結果ではなくまた状態を返すので、その結果を別のAggregatingMergeTreeに入れられます。元データをもう一度読まずに、サイトごとの状態から作ってください。サイトごとのユーザー数(数値)を足すと、同じ人を何度も数えます。