You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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聚合表,再通过物化视图异步聚合去重后的数据:
    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;
    
    注意:聚合的时效性依赖ReplacingMergeTree的合并时机,若需强实时性,可手动触发OPTIMIZE TABLE raw_deduplicated FINAL(不建议频繁执行,会增加集群负载)。

2. 改进插入阶段去重:状态表+窗口函数

针对跨块去重的问题,可通过维护去重状态表,在插入时实现全局去重:

  • 创建状态表
    存储已完成去重的ID,确保跨数据块的去重能力:
    CREATE TABLE deduplication_state
    (
        id String
    )
    ENGINE = ReplacingMergeTree()
    ORDER BY id;
    
  • 带去重逻辑的聚合写入
    通过窗口函数过滤最新数据,同时结合状态表排除已处理的ID:
    -- 写入聚合表
    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;
    
    该方案需确保状态表写入与聚合写入的原子性,可通过事务(ClickHouse 21.8+支持)或脚本批量执行避免重复处理。

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.19 13:30:30