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
相关产品推荐
相关产品推荐

