TT Lab
开始
学习 学习路径 课程

ClickHouse — 从内部理解列式分析数据库

分批写入汇总表,用求和与聚合状态得到与原始数据相同的结果

在 TT Lab 中继续学习

目标

把源事件的(site, day)汇总分别写入 SummingMergeTree 和 AggregatingMergeTree,并写出在合并之前也能得到与源数据完全一致答案的查询。确认无法用求和压缩的值必须存成状态,以及如何把这些状态再次合并成更大的单位。

为什么重要

汇总表能让仪表板变得便宜,但在合并结束之前,相同的键会分散在多行之中。在这种状态下,如果用 SELECT * 读取,或者把不同用户数存成数字再相加,会悄悄得出错误的答案。本实验通过停止合并来固定这种错误的状态,评分器会把你的查询结果与直接聚合源表 agg.hits 得到的值逐格核对。分了几次写入,则通过 system.part_log 的记录来确认。

步骤

  1. 创建数据库 agg 和表 agg.hits——列 ts DateTime, site LowCardinality(String), user_id UInt64, dur_ms UInt32, is_bot UInt8(按此顺序),引擎 MergeTree,ORDER BY (site, ts)。
  2. 运行一次 /opt/lab/fixtures/aggregating/hits.sql,写入 100 万行。
  3. 用 ENGINE = SummingMergeTree ORDER BY (site, day) 创建 agg.daily (site LowCardinality(String), day Date, hits UInt64, dur_total UInt64),并用 SYSTEM STOP MERGES agg.daily 停止合并,然后把 agg.hits 按时间段(toHour(ts) < 8、8 이상 16 미만、16 이상,后两个为韩文描述,分别表示 8 点(含)到 16 点之前、16 点(含)及以后)分三次 INSERT 写入 site, toDate(ts), count(), sum(dur_ms)。
  4. 创建在合并之前的现在,也能得出与源数据相同的按(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。
  5. 在 SYSTEM START MERGES agg.daily 之后,用 OPTIMIZE TABLE agg.daily FINAL 合并,使每个键成为一行。
  6. 用 ENGINE = AggregatingMergeTree ORDER BY (site, day) 创建 agg.daily_state (site LowCardinality(String), day Date, users AggregateFunction(uniqExact, UInt64), avg_dur AggregateFunction(avg, UInt32), human_hits AggregateFunction(countIf, UInt8)) 并停止合并,然后按与第 3 步相同的三个时间段,写入 uniqExactState(user_id)、avgState(dur_ms)、countIfState(is_bot = 0)。
  7. 创建只读取 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。
  8. 用 ENGINE = AggregatingMergeTree ORDER BY day 创建 agg.day_total (day Date, users AggregateFunction(uniqExact, UInt64)),把 agg.daily_state 的状态用 uniqExactMergeState(users) 按日期再次合并后写入,并创建只读取 agg.day_total、求出按日期的(不区分站点的)不同用户数的 /root/ch/aggregating/q_day.sql(列为 day, users)。

参考

源事件表

创建数据库 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,向 agg.hits 写入 1,000,000 行。

用 --queries-file 传给 clickhouse-client 即可。是 5 个站点 × 2026 年 9 月 30 日,所以(site, day)组合有 150 个。如果写了两次,就 TRUNCATE 之后重新写入。

分三次写入 SummingMergeTree

用 ENGINE = SummingMergeTree ORDER BY (site, day) 创建 agg.daily (site LowCardinality(String), day Date, hits UInt64, dur_total UInt64),并用 SYSTEM STOP MERGES agg.daily 停止合并,然后把 agg.hits 按时间段(toHour(ts) < 8、8 이상 16 미만、16 이상,后两个为韩文描述,分别表示 8 点(含)到 16 点之前、16 点(含)及以后)分三次 INSERT,把按(site, day)分组的 site, toDate(ts), count(), sum(dur_ms) 写入。

这是一天的汇总分三次到达的情形。每次 INSERT 只对那个时间段做 GROUP BY 后写入,所以合并之前,同一个(site, day)有三行。停止合并要在写入之前做。

合并之前也正确的查询

创建只读取 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 * 读取,现在每个键会有三行。不能读取源表 agg.hits——目的是只靠汇总表得出答案。

合并之后每个键一行

在 SYSTEM START MERGES agg.daily 之后,用 OPTIMIZE TABLE agg.daily FINAL 合并,使 agg.daily 成为一个数据片段、每个键一行。

合并会把相同键的数值列相加,折叠成一行。合并之后,看 SELECT * 是否与源数据的聚合一致。在生产环境中不知道合并的时机,所以像第 4 步那样的查询仍然是必要的。

无法用求和压缩的值存成状态

用 ENGINE = AggregatingMergeTree ORDER BY (site, day) 创建 agg.daily_state (site LowCardinality(String), day Date, users AggregateFunction(uniqExact, UInt64), avg_dur AggregateFunction(avg, UInt32), human_hits AggregateFunction(countIf, UInt8)),在 SYSTEM STOP MERGES agg.daily_state 之后,按与第 3 步相同的三个时间段,把 uniqExactState(user_id)、avgState(dur_ms)、countIfState(is_bot = 0) 按(site, day)分组后写入。

列类型 AggregateFunction(函数, 参数类型) 表示存放该函数的中间状态,写入时,在同一个函数上加 -State 来生成。-If 组合器把条件作为最后一个参数。因为是合并之前,所以三个批次的状态应当各自独立存在。

状态要先合并再计算结束

创建只读取 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 会就地把一个状态计算结束,所以把每行各自算出的数字相加,会把在多个时间段来过的人数多次。请看两个合计相差多少。

把状态再合并成更大的单位

用 ENGINE = AggregatingMergeTree ORDER BY day 创建 agg.day_total (day Date, users AggregateFunction(uniqExact, UInt64)),把 agg.daily_state 的 users 状态用 uniqExactMergeState(users) 按日期合并后写入,并创建只读取 agg.day_total、求出按日期(不区分站点)的不同用户数的 /root/ch/aggregating/q_day.sql(列为 day, users)。

-MergeState 会把状态合并之后,返回的不是结果,而是状态,所以可以把这个结果再放进另一个 AggregatingMergeTree。请不要重新读取源数据,而是从各站点的状态生成。如果把各站点的用户数(数字)相加,就会把同一个人数多次。