如何避免ClickHouse物化视图重复计算交易数据
问题描述
使用ClickHouse基于交易基表计算钱包余额时,若因异常导致基表重复插入相同交易,物化视图会重复计算这些交易,造成余额统计翻倍。需要找到避免物化视图重复统计已处理交易的方案。
现有表结构
交易基表
CREATE TABLE eth_trx_detail ( `hash` String, `blockHash` String, `blockNumber` Int32, `from` String, `to` String, `value` Decimal128(18), `time` Int32, `gasPrice` Int16, `fee` Decimal32(18) ) ENGINE = ReplacingMergeTree() ORDER BY (hash, blockNumber);
余额表
CREATE TABLE eth_address_balances ( `address` String, `balance` AggregateFunction(sum, Decimal(38, 18)), `last_updated` DateTime ) ENGINE = AggregatingMergeTree() ORDER BY (address);
物化视图
CREATE MATERIALIZED VIEW eth_balance_mv TO eth_address_balances AS SELECT address, sumState(value) as balance, max(time) as last_updated FROM ( -- 转出交易(余额减少) SELECT `from` as address, -value as value, time FROM eth_trx_detail UNION ALL -- 转入交易(余额增加) SELECT `to` as address, value as value, time FROM eth_trx_detail WHERE `to` != '' -- 排除合约创建交易 ) GROUP BY address;
问题重现
重复执行两次以下插入语句:
INSERT INTO blockbin_test.eth_trx_detail (*) VALUES (0xa5a9292698e6b2c3415d7bed272a5c86a2af398bf5b3f5c98cfa9e61b8f37c17,0x3fd256da659f74174bcc0c90f2bc0e7b5013914c94e819f24f0fe7c315bd606f,19840831,0x7591d15ca9c726faac98ef757f25009ce0efb1e9,0x8800d494d79b79955ec3575cc8bd400528853849,0.006,1715358527,75,0.0016134);
得到错误的余额结果:
┌─address──────────────┬─balance─┬────────last_updated─┐ │ 7.764412632850413e47 │ 0.012 │ 2024-05-10 16:28:47 │ │ 6.712037662395873e47 │ -0.012 │ 2024-05-10 16:28:47 │ └──────────────────────┴─────────┴─────────────────────┘
预期结果应为:
┌─address──────────────┬─balance─┬────────last_updated─┐ │ 7.764412632850413e47 │ 0.006 │ 2024-05-10 16:28:47 │ │ 6.712037662395873e47 │ -0.006 │ 2024-05-10 16:28:47 │ └──────────────────────┴─────────┴─────────────────────┘
解决方案
1. 利用ReplacingMergeTree的FINAL修饰符去重
交易基表使用ReplacingMergeTree,相同(hash, blockNumber)的记录会在合并时去重,但物化视图默认读取未合并的原始数据。在物化视图的查询中添加FINAL修饰符,强制读取合并后的去重数据:
CREATE MATERIALIZED VIEW eth_balance_mv TO eth_address_balances AS SELECT address, sumState(value) as balance, max(time) as last_updated FROM ( -- 转出交易(余额减少) SELECT `from` as address, -value as value, time FROM eth_trx_detail FINAL UNION ALL -- 转入交易(余额增加) SELECT `to` as address, value as value, time FROM eth_trx_detail FINAL WHERE `to` != '' -- 排除合约创建交易 ) GROUP BY address;
注意:FINAL会增加查询计算开销,适合数据量较小或合并频率较高的场景。
2. 新增处理状态标记过滤已处理交易
在交易基表中添加is_processed字段,标记交易是否已被物化视图处理:
ALTER TABLE eth_trx_detail ADD COLUMN is_processed UInt8 DEFAULT 0;
修改物化视图仅处理未标记的交易,处理完成后通过后台任务将is_processed更新为1:
CREATE MATERIALIZED VIEW eth_balance_mv TO eth_address_balances AS SELECT address, sumState(value) as balance, max(time) as last_updated FROM ( SELECT `from` as address, -value as value, time FROM eth_trx_detail WHERE is_processed = 0 UNION ALL SELECT `to` as address, value as value, time FROM eth_trx_detail WHERE is_processed = 0 AND `to` != '' ) GROUP BY address;
之后定期执行更新语句标记已处理交易:
ALTER TABLE eth_trx_detail UPDATE is_processed = 1 WHERE is_processed = 0;
这种方式适合需要精确控制处理逻辑的场景,但需额外维护后台任务。
3. 使用VersionedCollapsingMergeTree抵消重复交易
将交易基表改为VersionedCollapsingMergeTree,通过sign字段标记记录状态,重复插入的交易可插入抵消记录:
CREATE TABLE eth_trx_detail ( `hash` String, `blockHash` String, `blockNumber` Int32, `from` String, `to` String, `value` Decimal128(18), `time` Int32, `gasPrice` Int16, `fee` Decimal32(18), `sign` Int8 DEFAULT 1, `version` UInt64 DEFAULT now() ) ENGINE = VersionedCollapsingMergeTree(sign, version) ORDER BY (hash, blockNumber);
重复插入交易时,插入一条sign = -1的抵消记录:
INSERT INTO eth_trx_detail VALUES (..., -1, now());
合并后重复交易将被抵消,物化视图查询时自动处理这些记录,避免重复计算。
4. 先按唯一键去重再聚合
在物化视图的子查询中,先通过row_number()按交易唯一键(hash, blockNumber)去重,再进行余额计算:
CREATE MATERIALIZED VIEW eth_balance_mv TO eth_address_balances AS SELECT address, sumState(value) as balance, max(time) as last_updated FROM ( SELECT address, value * direction as value, time FROM ( -- 按交易唯一键去重,保留最新记录 SELECT hash, blockNumber, `from`, `to`, value, time, row_number() OVER (PARTITION BY hash, blockNumber ORDER BY time DESC) as rn FROM eth_trx_detail ) WHERE rn = 1 -- 转换出入账方向 UNPIVOT ( address, direction FOR type IN (`from` AS -1, `to` AS 1) ) WHERE NOT (direction = 1 AND address = '') -- 排除合约创建交易 ) GROUP BY address;
这种方式无需依赖MergeTree的合并逻辑,直接在查询层去重,适合对实时性要求较高的场景。
内容的提问来源于stack exchange,提问作者Amirhossein Masihi
相关产品推荐
相关产品推荐

