分批写入汇总表,用求和与聚合状态得到与原始数据相同的结果
目标
把源事件的(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,写入 100 万行。 - 用
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)。 - 创建在合并之前的现在,也能得出与源数据相同的按(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合并,使每个键成为一行。 - 用
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)。 - 创建只读取 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。 - 用
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)。
参考
- 服务器在 Pod 启动时就已经运行。只需输入
clickhouse-client即可连接。如果停了,就运行ch-up(重新启动服务器后 STOP MERGES 会被解除)。 - INSERT 记录会留在
system.part_log(event_type = 'NewPart')中。评分器按表的 UUID 过滤。想重来的话,先TRUNCATE TABLE,再写入三次即可。 - 状态列不是给人读的值。确认时,用
finalizeAggregation(열)(占位符为列名)或-Merge函数计算结束后再看。 - 常见错误:用一次 INSERT 写入全部源数据——相同的键在写入的那一刻就被合并,看不到“合并之前”。在查询中把
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,向 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。请不要重新读取源数据,而是从各站点的状态生成。如果把各站点的用户数(数字)相加,就会把同一个人数多次。