Clickhouse中如何实现Kafka数据流先去重再执行Rollup聚合?
ClickHouse先去重再聚合(Rollup)实现方案
问题描述
在数据处理场景中,去重与聚合(Rollup)是常见需求。ClickHouse虽支持ReplacingMergeTree(去重)和SummingMergeTree(聚合)两种引擎,但二者难以直接结合:无法通过物化视图将去重后的数据迁移至Rollup表,因为物化视图会在去重前的插入阶段触发(官方文档明确说明此限制)。
我们需要实现先按ID对Kafka流去重,再执行聚合的需求,已考虑过两种方案但存在缺陷:
- 插入阶段去重后写入聚合表:通过读取Kafka的物化视图,用
group by、distinct或row_number()窗口函数过滤rownum=1的方式完成去重,再写入SummingMergeTree表做聚合。但该方案仅能对单Kafka数据块内的数据去重,无法跨块去重,且去重窗口不可调整。 - 借助外部调度迁移数据:先用ReplacingMergeTree表完成去重,再通过外部定时调度器执行
INSERT INTO .. SELECT语句(配合FINAL或SQL去重逻辑)将数据迁移至SummingMergeTree表。但FINAL并不被官方推荐,且依赖外部调度组件,无法仅依托ClickHouse完成。
可行解决方案
1. 分层存储:ReplacingMergeTree + 异步聚合物化视图
这是最贴近官方设计思路的方案,依托ClickHouse自身的异步合并机制实现先去重再聚合,无需外部组件:
- 步骤1:创建去重表
用ReplacingMergeTree存储原始数据,按去重ID排序,指定版本字段(如事件时间)确保保留最新数据:CREATE TABLE raw_deduplicated ( id String, metric UInt64, event_time DateTime ) ENGINE = ReplacingMergeTree(event_time) ORDER BY id; - 步骤2:Kafka数据写入去重表
创建物化视图直接将Kafka流数据写入去重表,由ReplacingMergeTree后台自动完成异步去重:CREATE MATERIALIZED VIEW kafka_to_raw_deduplicated TO raw_deduplicated AS SELECT id, metric, event_time FROM kafka_topic; -- 需提前创建对应Kafka引擎表 - 步骤3:创建聚合表与物化视图
基于去重表创建SummingMergeTree聚合表,再通过物化视图异步聚合去重后的数据:
注意:聚合的时效性依赖ReplacingMergeTree的合并时机,若需强实时性,可手动触发CREATE TABLE aggregated_data ( id String, total_metric UInt64, dt Date ) ENGINE = SummingMergeTree(total_metric) ORDER BY (dt, id); CREATE MATERIALIZED VIEW mv_aggregated TO aggregated_data AS SELECT id, sum(metric) AS total_metric, toDate(event_time) AS dt FROM raw_deduplicated GROUP BY dt, id;OPTIMIZE TABLE raw_deduplicated FINAL(不建议频繁执行,会增加集群负载)。
2. 改进插入阶段去重:状态表+窗口函数
针对跨块去重的问题,可通过维护去重状态表,在插入时实现全局去重:
- 创建状态表
存储已完成去重的ID,确保跨数据块的去重能力:CREATE TABLE deduplication_state ( id String ) ENGINE = ReplacingMergeTree() ORDER BY id; - 带去重逻辑的聚合写入
通过窗口函数过滤最新数据,同时结合状态表排除已处理的ID:
该方案需确保状态表写入与聚合写入的原子性,可通过事务(ClickHouse 21.8+支持)或脚本批量执行避免重复处理。-- 写入聚合表 INSERT INTO aggregated_data SELECT id, sum(metric) AS total_metric, toDate(event_time) AS dt FROM ( SELECT id, metric, event_time, row_number() OVER (PARTITION BY id ORDER BY event_time DESC) AS rn FROM kafka_topic WHERE NOT exists (SELECT 1 FROM deduplication_state WHERE id = kafka_topic.id) ) WHERE rn = 1 GROUP BY dt, id; -- 更新状态表 INSERT INTO deduplication_state SELECT id FROM ( SELECT id, row_number() OVER (PARTITION BY id ORDER BY event_time DESC) AS rn FROM kafka_topic ) WHERE rn = 1;
3. 实时查询聚合(无需持久化聚合结果)
若无需持久化聚合数据,可直接在查询去重表时实时聚合,适合查询频率低、数据量不大的场景:
SELECT toDate(event_time) AS dt, id, sum(metric) AS total_metric FROM raw_deduplicated GROUP BY dt, id
若需确保查询时已完成去重,可添加FINAL关键字,但会显著降低查询性能,仅建议在小数据集场景使用:
SELECT toDate(event_time) AS dt, id, sum(metric) AS total_metric FROM raw_deduplicated FINAL GROUP BY dt, id
内容的提问来源于stack exchange,提问作者kev
相关产品推荐
相关产品推荐

