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

ClickHouse基于多列规则的去重/Upsert实现方案问询

实现ClickHouse基于(id, property_id)的Upsert逻辑

基础方案(不含JSON合并)

通过「源表 + AggregatingMergeTree目标表 + 物化视图」的流水线实现核心Upsert规则,解决你之前遇到的主键、分区、聚合函数问题:

1. 创建聚合目标表

用AggregatingMergeTree存储聚合状态,排序键设置为去重依据(id, property_id),确保同一组数据能被正确聚合:

CREATE TABLE test_upsert_target
(
    id Int,
    property_id Int,
    created_at_min AggregateFunction(min, DateTime('UTC')),
    modified_at_max AggregateFunction(max, DateTime('UTC')),
    data_latest AggregateFunction(argMax, Int, DateTime('UTC'))
) ENGINE = AggregatingMergeTree()
PARTITION BY toYYYYMM(minMerge(created_at_min))  -- 基于最早的created_at分区,保证同组数据落在同一分区
ORDER BY (id, property_id);

2. 创建同步物化视图

将源表数据实时聚合到目标表,使用*State系列函数生成聚合状态:

CREATE MATERIALIZED VIEW test_upsert_mv TO test_upsert_target
AS
SELECT
    id,
    property_id,
    minState(created_at) AS created_at_min,
    maxState(modified_at) AS modified_at_max,
    argMaxState(data, modified_at) AS data_latest
FROM test
GROUP BY id, property_id;

3. 数据写入与查询

  • 写入:直接向源表test插入数据,物化视图会自动同步聚合状态到目标表
  • 查询:用*Merge系列函数解析聚合状态,得到最终去重结果:
SELECT
    id,
    property_id,
    minMerge(created_at_min) AS created_at,
    maxMerge(modified_at_max) AS modified_at,
    argMaxMerge(data_latest) AS data
FROM test_upsert_target
GROUP BY id, property_id;

历史数据初始化

如果源表已有历史数据,需手动同步一次:

INSERT INTO test_upsert_target
SELECT
    id,
    property_id,
    minState(created_at),
    maxState(modified_at),
    argMaxState(data, modified_at)
FROM test
GROUP BY id, property_id;

进阶方案(含JSON对象合并)

ClickHouse无内置JSON合并聚合函数,需通过「JSON解析为Map + 聚合Map + 转回JSON」实现,以下是两种常见场景:

场景1:合并所有JSON键(重复键保留最后插入值)

1. 修改目标表与物化视图

-- 新增Map类型的JSON聚合字段
CREATE TABLE test_upsert_target_with_json
(
    id Int,
    property_id Int,
    created_at_min AggregateFunction(min, DateTime('UTC')),
    modified_at_max AggregateFunction(max, DateTime('UTC')),
    data_latest AggregateFunction(argMax, Int, DateTime('UTC')),
    json_merged AggregateFunction(mergeMaps, Map(String, String))
) ENGINE = AggregatingMergeTree()
PARTITION BY toYYYYMM(minMerge(created_at_min))
ORDER BY (id, property_id);

-- 物化视图解析JSON为Map并聚合
CREATE MATERIALIZED VIEW test_upsert_mv_with_json TO test_upsert_target_with_json
AS
SELECT
    id,
    property_id,
    minState(created_at) AS created_at_min,
    maxState(modified_at) AS modified_at_max,
    argMaxState(data, modified_at) AS data_latest,
    mergeMapsState(
        if(json_str IS NOT NULL, JSONExtractKeysAndValues(json_str, 'String'), map())
    ) AS json_merged
FROM test
GROUP BY id, property_id;

2. 查询转回JSON字符串

SELECT
    id,
    property_id,
    minMerge(created_at_min) AS created_at,
    maxMerge(modified_at_max) AS modified_at,
    argMaxMerge(data_latest) AS data,
    mapToString(mergeMapsMerge(json_merged)) AS json_str
FROM test_upsert_target_with_json
GROUP BY id, property_id;

场景2:重复JSON键保留modified_at最新值

需拆分JSON为键值对后单独聚合,再重组为JSON:

-- 1. 拆分JSON为键值对的物化视图
CREATE MATERIALIZED VIEW test_json_kv_mv AS
SELECT
    id,
    property_id,
    modified_at,
    key,
    value
FROM test
ARRAY JOIN JSONExtractKeysAndValues(json_str, 'String') AS (key, value)
WHERE json_str IS NOT NULL;

-- 2. 聚合每个键的最新值
CREATE TABLE test_json_agg
(
    id Int,
    property_id Int,
    key String,
    latest_value AggregateFunction(argMax, String, DateTime('UTC'))
) ENGINE = AggregatingMergeTree()
ORDER BY (id, property_id, key);

CREATE MATERIALIZED VIEW test_json_agg_mv TO test_json_agg AS
SELECT
    id,
    property_id,
    key,
    argMaxState(value, modified_at) AS latest_value
FROM test_json_kv_mv
GROUP BY id, property_id, key;

-- 3. 最终查询重组JSON
SELECT
    t.id,
    t.property_id,
    t.created_at,
    t.modified_at,
    t.data,
    mapToString(groupMap(key, argMaxMerge(latest_value))) AS json_str
FROM (
    SELECT
        id,
        property_id,
        minMerge(created_at_min) AS created_at,
        maxMerge(modified_at_max) AS modified_at,
        argMaxMerge(data_latest) AS data
    FROM test_upsert_target
    GROUP BY id, property_id
) t
LEFT JOIN test_json_agg j ON t.id = j.id AND t.property_id = j.property_id
GROUP BY t.id, t.property_id, t.created_at, t.modified_at, t.data;

常见问题说明

  • 主键/排序键:目标表ORDER BY必须设置为去重键(id, property_id),否则无法触发正确聚合
  • 分区键:使用聚合后的固定值(如最早的created_at),避免同一组数据跨分区导致聚合不完整
  • 聚合函数:物化视图用*State生成聚合状态,查询用*Merge解析状态,二者需一一对应

内容的提问来源于stack exchange,提问作者aimfeld

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.11 18:35:23