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

ClickHouse可刷新物化视图导致分布式查询结果不一致问题求助

ClickHouse可刷新物化视图导致分布式表数据不一致问题

问题概述

搭建了带有可刷新物化视图的ClickHouse分布式集群(版本25.8.8.26-lte),查询目标分布式表some_db.result时无法获取全部数据,结果随时间不一致:每隔一分钟查询,数据仅来自某一个分片的local_result表,其他分片该表为空;刷新周期过后,数据切换到另一个分片,原分片数据被清空。

集群配置SQL

CREATE DATABASE IF NOT EXISTS some_db ON CLUSTER some_cluster
ENGINE = Replicated('/clickhouse/databases/some_db', '{shard}', '{replica}');

-- 源表定义
CREATE TABLE some_db.local_source
(
  id String CODEC(ZSTD(1)),
  a  Int8 DEFAULT 0 CODEC(T64, ZSTD(1)),
  b  Int8 DEFAULT 0 CODEC(T64, ZSTD(1)),
)
ENGINE = ReplicatedMergeTree
ORDER BY id;

CREATE TABLE some_db.source as some_db.local_source
ENGINE = Distributed(some_cluster, some_db, local_source, metroHash64(id));

-- 聚合中间表定义
CREATE TABLE some_db.local_calculated
(
  id           String CODEC(ZSTD(1)),
  a_sum_state  AggregateFunction(sum, Int64) CODEC(ZSTD(1)),
  b_sum_state  AggregateFunction(sum, Int64) CODEC(ZSTD(1)),
)
ENGINE = ReplicatedAggregatingMergeTree
ORDER BY id;

CREATE TABLE some_db.calculated as some_db.local_calculated
    ENGINE = Distributed(some_cluster, some_db, local_calculated, metroHash64(id));

CREATE MATERIALIZED VIEW some_db.mv_calc_characteristics
TO some_db.local_calculated
AS SELECT
    s.id                   AS id,
    sumState(toInt64(s.a)) AS a_sum_state,
    sumState(toInt64(s.b)) AS b_sum_state
FROM some_db.local_source s
GROUP BY id;

-- 结果表与可刷新物化视图定义
CREATE TABLE some_db.local_result
(
  id      String CODEC(ZSTD(1)),
  total_a Int64  CODEC(T64, ZSTD(1)),
  total_b Int64  CODEC(T64, ZSTD(1)),
  is_c    Int64  CODEC(T64, ZSTD(1)),
)
ENGINE = ReplicatedMergeTree
ORDER BY id;

CREATE TABLE some_db.result as some_db.local_result
ENGINE = Distributed(some_cluster, some_db, local_result, metroHash64(id));

CREATE MATERIALIZED VIEW some_db.mv_calc_categories
REFRESH EVERY 1 MINUTE TO some_db.local_result
AS SELECT
    id                                   AS id,
    sumMerge(y.a_sum_state)              AS total_a,
    sumMerge(y.b_sum_state)              AS total_b,
    (total_a > 0 AND total_b > 0)        AS is_c
FROM some_db.local_calculated y
GROUP BY id;

数据插入与查询SQL

插入数据

INSERT INTO some_db.source SELECT
 (100 + (rand() % (900 - 100 + 1))) as id,
 randBernoulli(0.3) AS a,
 randBernoulli(0.3) AS b
FROM numbers(10);

查询结果表

SELECT
    id,
    total_a,
    total_b,
    is_c
FROM some_db.result LIMIT 10;

问题原因

  1. 可刷新物化视图默认替换模式:带REFRESH EVERY的物化视图默认使用REFRESH MODE REPLACE,每次刷新会先清空目标表(local_result),再插入本次计算结果。
  2. 本地表计算范围限制:每个分片的可刷新MV仅查询当前分片的local_calculated表,导致每个分片的local_result仅包含自身分片的聚合数据。
  3. 分布式表汇总逻辑冲突:查询分布式表result时会汇总所有分片的local_result数据,但由于每个MV刷新时会清空自身分片的local_result,导致不同时间点只有刚完成刷新的分片有数据,其他分片为空,最终结果随刷新周期交替变化。

解决方案

方案1:修改可刷新MV为追加模式

将可刷新物化视图的刷新模式改为APPEND,避免每次刷新清空目标表,同时将目标表改为ReplacingMergeTree处理重复数据:

-- 修改目标表引擎为ReplacingMergeTree,用刷新时间作为版本字段
ALTER TABLE some_db.local_result ON CLUSTER some_cluster
ENGINE = ReplicatedReplacingMergeTree
ORDER BY id
SETTINGS index_granularity = 8192;

-- 重新创建可刷新MV,指定APPEND模式并添加版本字段
DROP MATERIALIZED VIEW IF EXISTS some_db.mv_calc_categories ON CLUSTER some_cluster;
CREATE MATERIALIZED VIEW some_db.mv_calc_categories ON CLUSTER some_cluster
REFRESH EVERY 1 MINUTE
REFRESH MODE APPEND
TO some_db.local_result
AS SELECT
    id                                   AS id,
    sumMerge(y.a_sum_state)              AS total_a,
    sumMerge(y.b_sum_state)              AS total_b,
    (total_a > 0 AND total_b > 0)        AS is_c,
    now()                                AS version -- 用于ReplacingMergeTree去重
FROM some_db.local_calculated y
GROUP BY id;

查询时通过FINAL关键字获取去重后的最新数据:

SELECT
    id,
    total_a,
    total_b,
    is_c
FROM some_db.result FINAL LIMIT 10;

方案2:让可刷新MV查询分布式聚合表

让每个分片的可刷新MV查询全量的分布式聚合表some_db.calculated,确保每个分片的local_result包含全量数据:

DROP MATERIALIZED VIEW IF EXISTS some_db.mv_calc_categories ON CLUSTER some_cluster;
CREATE MATERIALIZED VIEW some_db.mv_calc_categories ON CLUSTER some_cluster
REFRESH EVERY 1 MINUTE
TO some_db.local_result
AS SELECT
    id                                   AS id,
    sumMerge(y.a_sum_state)              AS total_a,
    sumMerge(y.b_sum_state)              AS total_b,
    (total_a > 0 AND total_b > 0)        AS is_c
FROM some_db.calculated y -- 改为查询分布式表
GROUP BY id;

此方案下每个分片的local_result都会包含全量数据,查询分布式表时不会出现数据缺失,但会增加集群计算压力,适合小数据量场景。

方案3:普通物化视图结合定时任务

如果不需要严格分钟级刷新,改用普通物化视图,通过外部定时任务执行刷新命令:

-- 创建普通物化视图
DROP MATERIALIZED VIEW IF EXISTS some_db.mv_calc_categories ON CLUSTER some_cluster;
CREATE MATERIALIZED VIEW some_db.mv_calc_categories ON CLUSTER some_cluster
TO some_db.local_result
AS SELECT
    id                                   AS id,
    sumMerge(y.a_sum_state)              AS total_a,
    sumMerge(y.b_sum_state)              AS total_b,
    (total_a > 0 AND total_b > 0)        AS is_c
FROM some_db.calculated y
GROUP BY id;

在集群节点上定时执行刷新命令:

# 每分钟刷新一次物化视图(需确保所有分片执行,或通过分布式命令)
clickhouse-client -q "REFRESH MATERIALIZED VIEW some_db.mv_calc_categories ON CLUSTER some_cluster"

内容的提问来源于stack exchange,提问作者Егор Лебедев

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.12 06:44:51