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

如何避免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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.12 14:25:56