帮你快速理解、总结文档立即下载
文档中心>腾讯云数据仓库 TCHouse-C>开发指南>表引擎>自研选择性复制表 Selective Replication

自研选择性复制表 Selective Replication

最近更新时间:2026-09-18 15:52:30
我的收藏
腾讯云数据仓库 TCHouse-C 自研了选择性复制表(Selective Replication,简称 SR),作为对 ClickHouse 副本表 ReplicatedMergeTree 的重要扩展能力。本文为您介绍选择性复制表的工作原理、与 ClickHouse 原生分片副本方案的差异、建表配置、读写方式、容量规划方法、运维操作与故障处理,以及典型应用场景。
说明:
当前选择性复制表处于白名单试用阶段、仅在 弹性版 中支持。若您想要体验该功能,可 联系我们 为您开通白名单。

功能介绍

ClickHouse 的副本(Replication)与分片(Sharding)是其实现高可用与水平扩展的两大核心机制,副本解决数据可靠性与读负载均衡,分片解决单机存储上限与写/计算吞吐瓶颈。其中:
分片:会将一张大表的数据按某种策略(如哈希、随机或按列值)切分成多个数据块,分散存储在不同的节点上,每个分片只保存全量数据的一部分,通过 Distributed 表引擎来实现逻辑聚合。查询时,Distributed 表将请求分发至各个分片并行计算,再汇总结果返回。
副本:针对同一个分片的数据,在不同节点上保存互为镜像的物理副本。依赖 ReplicatedMergeTree 表引擎系列,借助 ZooKeeper 或 ClickHouse Keeper 在副本间采用异步多主(Async Multi-Master)机制进行数据同步。
这一实现机制在以下场景中会存在局限性:
新增分片节点后,已存在的数据不会自动重新分布到新节点上,需手动操作转移 Data Part 或重新写入,导致集群弹性扩缩容成本极高;
在多分片架构下,跨节点的分布式 JOIN 需要在网络间传输大量数据,在大表关联场景下性能低于本地 JOIN;
集群拓扑由节点决定、集群上所有表共用「N 个分片 × M 个副本」布局,对可用性/数据可靠性要求低的表依然要强制占用 M 份副本的存储资源。
选择性复制表在保持 ReplicatedMergeTree 副本语义的基础上,将分片数与副本数从集群配置下移到表 DDL,表的每个 Partition 仅存储在集群中 replication_factor 个节点而非全部节点、可按表的重要程度与数据量分别设置副本数,数据布局不再受集群拓扑约束、节点增减无需变更配置或手工迁移数据,因此有以下收益:
收益
具体表现
收益
扩缩容自动重分布
节点加入集群后自动 Rebalance,存量 Partition 自动迁移。
扩容无需变更集群配置、无需手工搬迁数据,也无需重建表。
大表关联查询性能提升
分桶配置一致的两张表自动本地 JOIN。
无需将右表广播到全部节点,跨节点传输量随之下降,大表关联不再受右表体积限制。
布局按表定义
分桶数与副本数是表属性,不再是集群属性。
同一集群上核心表与日志表可分别设置副本数,调整时影响范围收敛到单张表、不波及集群上的其他业务。
降低存储成本
每个 Partition 仅存储 replication_factor 份,而非节点数份。
N 节点集群设置 replication_factor = 2 时,单表存储占用为原生方案的 2/N。以 10 节点集群为例,存储占用由 10 份降至 2 份。
说明:
存储收益与可靠性是一组权衡,replication_factor = 2 时每个 Partition 仅有 2 份数据,只能容忍 1 个节点故障;如需容忍 2 个节点同时故障,须设置 replication_factor 不小于 3。
选择性复制表与社区原生 ReplicatedMergeTree 的对比:
维度
原生方案(Distributed + ReplicatedMergeTree)
选择性复制表(ReplicatedMergeTree)
分片与副本的声明位置
集群配置统一声明
表 DDL 的 SHARD BY ... INTO N 与 replication_factor
生效粒度
节点。节点在拓扑中的位置决定其持有的数据
表。同一节点为不同表承载不同份额
同集群多表的布局
强制统一,所有表继承同一套分片副本结构
各表独立,可分别设置分桶数与副本数
分片数的来源
集群配置中声明的分片数,调整需变更集群配置
表级 INTO N,与集群节点数无关,取值范围 2 至 256
副本数的来源
分片内配置的节点数,集群上所有表相同
表级 replication_factor,可按表设置
节点间的数据关系
同分片内各节点数据完全一致,跨分片无交集
每个 Partition 由 replication_factor 个节点持有,任意两节点的数据集部分重叠
数据落点的决定方式
由节点在拓扑中的身份静态确定
由引擎按 Partition 计算,结果记录于 Keeper
路由实现层
Distributed 表层,需显式建表并指定 sharding_key
存储引擎层,对客户端透明
用户可见表数量
2 张,副本表与分布式表
1 张
集群扩容
需修改集群配置,存量数据保持原分布,不自动迁移
节点加入即可,触发 rebalance 后存量 partition 自动重分布
布局调整的影响范围
集群级,改动波及该集群上所有表
表级,重建单表即可,其余表不受影响
本地 JOIN
两表使用同一 sharding_key 且分片布局一致,需要自行保证
分桶配置一致时引擎自动本地执行

典型使用场景推荐

当前推荐用户使用的典型场景如下,完整的建表语句与操作示例请参见后续文档。
场景
业务诉求
关键配置
选择性复制表收益
大表关联分析
两张亿级明细表按同一业务键关联,广播右表已成为瓶颈
SHARD BY user_id INTO N,两表的 N 与列类型必须一致
两表按同一键分桶后在节点本地完成关联,不广播右表
存储成本优化
集群存储水位偏高,扩容受成本约束
replication_factor = 2,须先评估可接受的故障容忍度
副本份数与节点数解耦,由 replication_factor 决定,单机存储占用降至约 2 ÷ 节点数
冷热表差异化副本策略
核心交易表要求高可靠,日志表以省空间为先
核心表 replication_factor = 3,日志表 replication_factor = 2
副本数按表设置,同集群可并存多种策略
弹性扩缩容
大促前后增减节点,不希望停机搬迁数据
控制单个 partition 体积,避免迁移触达超时上限
节点水平增减后,自动重分布存量 Partition
明细日志高速写入
只按时间范围扫描,不做键关联,追求写入吞吐
SHARD RANDOM INTO N,须自行处理写入幂等
随机分桶将写入均摊到各节点、避免热点,写入吞吐最高

使用说明

已支持的操作

以下操作在选择性复制表上可用:
类别
操作
说明
数据变更
Mutation(ALTER DELETE、ALTER UPDATE)、轻量删除(DELETE FROM)
支持,在各归属副本上分别执行
生命周期
TTL(表级、列级、TTL MOVE)
支持
查询
OPTIMIZE FINAL、SELECT ... FINAL
支持。使用合并语义引擎时,去重结果的正确性仍取决于分桶键是否被 ORDER BY 覆盖,请参见 合并语义引擎的使用约束。
加速结构
跳数索引、投影、物化视图
支持
表结构
ALTER ADD / DROP / MODIFY COLUMN
支持,_shard 列除外
分区操作
DROP PARTITION ID
支持,须逐分桶执行。
查询优化
分区裁剪
支持
合并
OPTIMIZE
支持,自动转发至归属副本
副本同步
SYSTEM SYNC REPLICA
支持
写入一致性
insert_quorum
支持,受 replication_factor 上限约束
说明:
上述能力在 选择性复制表上的执行范围是该 partition 的归属副本,而非集群全部节点。因此 mutation 与 TTL 的执行进度需按 partition 观察,不能以单个节点的完成情况判断全表状态。查询 mutation 进度请使用 system.mutations,并结合 系统表 确认各 partition 的归属副本。

不支持的操作

在使用选择性复制表时请确认以下限制,当前不支持的操作有:
类别
操作
说明
数据布局变更
ALTER ... MODIFY SETTING replication_factor
不支持。请改用专用 DDL ALTER TABLE ... MODIFY REPLICATION FACTOR N
数据布局变更
RESET SETTING replication_factor
不支持
数据布局变更
修改分桶键或 INTO N
不支持,无对应语法。只能新建表并迁移数据
分区搬迁
ALTER ... DETACH PARTITION
不支持
分区搬迁
ALTER ... ATTACH PARTITION
不支持。ATTACH PART 可用
备份恢复
BACKUP、RESTORE、FETCH PARTITION、FREEZE PARTITION
不支持。选择性复制表的 Partition 分散在部分节点上,且归属关系记录于 Keeper,上述命令的单机语义无法覆盖
组合使用
在选择性复制表之上创建 Distributed 表
不支持。选择性复制表的路由在存储引擎层完成,无需也不应再叠加 Distributed 层
列操作
修改 _shard 列
不支持,该列由 SHARD BY 子句维护
写入一致性
insert_quorum 大于 replication_factor
不支持,选择性复制表仅由归属副本确认写入
说明:
若业务有明确的备份恢复要求,请评估以下替代方案后再决定是否启用选择性复制表:
通过 INSERT INTO ... SELECT 将数据导出至对象存储或另一张表;
对关键表设置 replication_factor = 0 使用标准 ReplicatedMergeTree 以保留原生备份能力;
提高 replication_factor 以增强在线可靠性、降低对离线备份的依赖。
注意:
使用合并语义引擎(ReplacingMergeTree、SummingMergeTree 等)且分桶键未被 ORDER BY 完全覆盖。当前版本不会拦截该配置,但会导致去重语义失效。

前提条件

已创建 TCHouse-C 弹性版集群(存算分离架构),且集群节点数大于计划设置的 replication_factor。节点数等于 replication_factor 时,SR 退化为全节点存储全量数据、无存储收益。
已获取具备目标数据库建表和读写权限的数据库账号。
已通过客户端或控制台 SQL 工作区连接到集群。

快速入门

以下示例完成建表、写入、查询与归属关系确认。前提条件为 3 节点集群,各节点的副本名分别为 r1、r2、r3。

步骤 1:建表

在每个节点上分别执行建表语句。
CREATE TABLE events
(
event_date Date,
user_id UInt64,
action String
)
ENGINE = ReplicatedMergeTree
PARTITION BY toYYYYMM(event_date)
SHARD BY cityHash64(user_id) INTO 16
ORDER BY (user_id, event_date)
SETTINGS replication_factor = 2;

步骤 2:写入与查询

写入可发往任意节点,引擎自动转发至归属副本。
INSERT INTO events VALUES
('2024-01-15', 1001, 'login'),
('2024-01-15', 1002, 'logout'),
('2024-02-01', 2001, 'purchase');

SELECT count() FROM events;
SELECT * FROM events WHERE user_id = 1001;

步骤 3:确认分区归属关系

SELECT partition_id, assigned_replicas, is_under_replicated
FROM system.selective_assignments
WHERE database = currentDatabase() AND table = 'events'
ORDER BY partition_id;
-- 预期结果:partition_id 形如 202401-0 至 202401-15
-- 每行 assigned_replicas 包含 2 个副本,is_under_replicated 为 0

建表

本节介绍建表基本语法、分桶模式选择、参数含义与校验规则。

基本语法

CREATE TABLE [db.]table_name
(
col1 Type1,
...
)
ENGINE = ReplicatedMergeTree
[PARTITION BY <user_partition_expr>]
[ SHARD BY <shard_expr> INTO <N> | SHARD RANDOM INTO <N> ]
ORDER BY <order_expr>
SETTINGS replication_factor = <M>;
约束如下:
SHARD BY 子句可选。弹性版集群默认开启选择性副本表能力,不配置 SHARD BY 的表同样是选择性复制表。
INTO <N> 在两种形式下均为必填,取值须为整数字面量。
存储子句顺序不限,但 SETTINGS 必须位于最后。
SHOW CREATE TABLE 的输出顺序固定为 ENGINE、PARTITION BY、SHARD BY、PRIMARY KEY、ORDER BY、SAMPLE BY、TTL、SETTINGS。

分桶模式

SHARD BY <expr> INTO <N> -- 确定性分桶,桶号为 positiveModulo(<expr>, N)
SHARD RANDOM INTO <N> -- 随机分桶,桶号为 modulo(randConstant(), N)
两种模式的能力差异:
特性
SHARD BY <expr> INTO N
SHARD RANDOM INTO N
写入重试幂等
支持
不支持
单 Block 拆分 Part 数
最多 N 个
1 个
本地 JOIN
支持
不支持
合并语义引擎
需满足分桶键约束
不应使用
SHARD RANDOM 的桶号由 randConstant() 生成,每个 Block 取值一次,因此单个 INSERT Block 仅落入一个分桶,不产生写放大。代价是重试时桶号不同,无法完成去重,且不参与本地 JOIN 优化。
说明:
选择分桶模式时,先确认业务是否需要写入幂等、合并语义或本地 JOIN。这三项能力中任意一项为必需时,都只能使用确定性 SHARD BY。SHARD RANDOM 仅适用于纯明细追加写入且不做大表关联的场景。

参数说明

表级参数在 CREATE TABLE ... SETTINGS 中配置:
参数
类型
默认值
说明
replication_factor
UInt64
2
每个 Partition 的副本数。取 0 表示关闭选择性复制表能力,与开源 ReplicatedMergeTree 行为保持一致。
max_concurrent_partition_migrations
UInt64
16
并发迁移的 Partition 数上限
migration_timeout_seconds
UInt64
300
迁移 CLONE 阶段超时时间,超时后回滚,单位秒
migration_coordinator_timeout
UInt64
60
迁移协调者失联判定时间,单位秒
migration_max_switch_delay
UInt64
60
SWITCH 阶段最大延迟,单位秒
enable_auto_rebalance
Bool
false
是否自动触发 Rebalance。默认开启
migration_monitor_interval_seconds
UInt64
10
后台迁移监控周期,单位秒
查询级参数:
参数
默认值
说明
selective_replication_auto_global_in_and_join
true
不符合分布规则时自动改写为 GLOBAL JOIN 或 GLOBAL IN
selective_replication_route_by_delay
true
按副本复制延迟择优路由
selective_replication_weighted_routing
true
按实际数据量加权路由
skip_unavailable_shards
false
是否容忍无可用副本的 Partition
max_number_of_shards
256
单个分区的最大分桶数
random_shard_insert_mode
prefer_local
随机分桶表的写入路由策略

建表校验规则

分桶数取值范围为 [2, 256],越界时报错:
SHARD BY INTO N requires N >= 2, got 1
SHARD BY INTO N requires N <= 256, got 512
分桶表达式须为单个整数类型表达式:
SHARD BY expects a single expression, got a tuple of 2
SHARD BY expression must have an integer type, but `name` has type String.
Wrap it into a hash function such as `cityHash64`.
_shard 为保留列名,不可声明或修改:
Column `_shard` is reserved for the SHARD BY clause and cannot be defined by the user
Cannot ALTER column `_shard` because it is maintained by the SHARD BY clause
使用确定性 SHARD BY 时,PARTITION BY 的每一列须为整数、Date、DateTime 或 IPv4 类型。例如 PARTITION BY country(String 类型)将被拒绝,需改为 PARTITION BY cityHash64(country)。该限制不适用于 SHARD RANDOM。
PARTITION BY column #1 has non-integral type 'String' — SHARD BY tables with role_aware
require every PARTITION BY column to be an integral, Date/DateTime, or IPv4 type.
Wrap the expression with cityHash64() or similar to obtain an integer.
在后续节点上建表时,内核会比对 Keeper 中已记录的 SR 配置,不一致时报错:
replication_factor mismatch: table in ZooKeeper has replication_factor=2, but this replica
has replication_factor=3. All replicas must use the same value

First replica was created without selective replication (replication_factor=0), but this
replica has replication_factor=2. All replicas must use the same value
说明:
服务启动阶段检测到配置不一致时,不会向客户端返回错误,而是将该节点置为只读。这种情况下写入会失败但没有明确的建表报错。

合并语义引擎的使用约束

ReplacingMergeTree、CollapsingMergeTree、VersionedCollapsingMergeTree、SummingMergeTree、AggregatingMergeTree、GraphiteMergeTree 的去重、折叠、聚合语义仅在同一 partition_id 内生效。分桶会将业务分区内 ORDER BY key 相同的行拆分到不同的物理 Partition,导致上述语义失效。
因此使用合并语义引擎时须同时满足:
采用确定性 SHARD BY <expr>,且 <expr> 引用的所有列均包含在 ORDER BY 中。
不使用 SHARD RANDOM。
正确与错误配置示例:
-- 正确:SHARD BY 的 user_id 包含在 ORDER BY 中
CREATE TABLE ok (event_date Date, user_id UInt64, v UInt64)
ENGINE = ReplicatedReplacingMergeTree
PARTITION BY toYYYYMM(event_date)
SHARD BY user_id INTO 8
ORDER BY (user_id, event_date)
SETTINGS replication_factor = 2;

-- 错误:v 未包含在 ORDER BY 中,同 key 数据被拆至不同 partition,去重失效
-- 当前版本不会报错,请勿使用该配置
CREATE TABLE bad (event_date Date, user_id UInt64, v UInt64)
ENGINE = ReplicatedReplacingMergeTree
PARTITION BY toYYYYMM(event_date)
SHARD BY v INTO 8
ORDER BY (user_id, event_date)
SETTINGS replication_factor = 2;
UniqueMergeTree 系列的对应校验已生效,违反时报错:
SHARD RANDOM cannot be used with ReplicatedUniqueMergeTree because random sharding breaks
deduplication semantics

UniqueMergeTree with SHARD BY requires all shard key columns to be part of the unique key.
Column 'uid' is used in SHARD BY but is not in the unique key [id]
说明:
如已知问题所述,当前版本对上述两条约束的建表校验未生效。使用合并语义引擎时,请在执行 CREATE TABLE 前人工核对 SHARD BY 引用的每一列是否都出现在 ORDER BY 中。违反约束不会报错,但去重会静默失效,数据长期堆积后才会被发现。

不配置 SHARD BY 的选择性复制表

不配置 SHARD BY 的表同样启用选择性复制表,这是默认形态:
CREATE TABLE t (key UInt64, v String)
ENGINE = ReplicatedMergeTree
ORDER BY key;
两种形态的差异:
维度
配置确定性 SHARD BY
未配置 SHARD BY
分配算法
Rendezvous 哈希,结果可预测
最少负载优先,结果依赖当前负载
本地 JOIN
支持
不支持
partition_id 格式
<业务分区>-<桶号>
<业务分区>

partition_id 格式

建表配置
partition_id 示例
PARTITION BY toYYYYMM(date) SHARD BY user_id INTO 8
202401-3
SHARD BY user_id INTO 8,无 PARTITION BY
3
PARTITION BY toYYYYMM(date),无 SHARD BY
202401
无 PARTITION BY 且无 SHARD BY
all
partition_id 格式变化会影响所有按 Partition 操作的语句。

角色配置

选择性副本表按节点发布的角色分组,将 replication_factor 份数据按配额分摊到各角色。角色在服务端配置文件中声明,不是表级参数:
<clickhouse>
<role>read</role>
</clickhouse>
角色名不可为空,且不可以下划线开头,下划线前缀为内核保留。
所有节点使用同一角色时,配额分摊退化为全节点平摊,等价于纯哈希分配。
角色对读取的影响为亲和性而非隔离:同角色的归属副本优先被选中,跨角色节点仍作为备选。
说明:
角色不区分读写,任一角色均会分配数据并承载写入。若您的目标是读写资源隔离,角色配置无法达成该效果,请使用 多计算资源组 能力。

数据写入

写入可发往任意节点,引擎按 Partition 归属关系转发至对应节点,无需在客户端实现路由。
INSERT INTO events VALUES
('2024-01-15', 1001, 'login'),
('2024-01-15', 1002, 'logout');

随机分桶表的写入路由

表使用 SHARD RANDOM 时,random_shard_insert_mode 控制写入落点:
取值
行为
prefer_local(默认)
优先写入接收节点;接收节点没有对应分区时,跨节点写入其他合适节点
local_only
仅写入接收节点;接收节点没有对应分区时写入失败,错误类型为 NO_AVAILABLE_REPLICA
random
随机写入合适节点

批量写入

单个 INSERT Block 最多拆分为 P_batch × N 个 Part,其中 P_batch 为本批数据跨越的业务分区数,N 为分桶数。分桶数较大时容易触达单次写入的分区数上限,需调整参数:
SET max_partitions_per_insert_block = 1000;
使用 SHARD RANDOM 可规避该问题,代价参见分桶模式。

写入失败的错误信息

无可用归属副本:
Selective replication: no replicas assigned for partition 202401-3 (original: 202401-3).
Check cluster health and rebalance state.
全部归属副本写入失败:
Selective replication INSERT: all assigned replicas [r1, r2] failed for 3 partition(s).
Last error: ...

insert_quorum 约束

选择性复制表仅由 Partition 的归属副本确认写入,因此 insert_quorum 不得大于 replication_factor:
insert_quorum (3) cannot be greater than replication_factor (2) in selective replication mode.
Only the 2 assigned replicas for each partition will confirm the write.
Lower insert_quorum to at most 2.
可用归属副本数不满足 Quorum 要求时:
Selective replication: partition 202401-3 has only 1 effective assigned replica(s), but
insert_quorum requires 2. Either reduce insert_quorum or wait for more replicas to become
available.

数据查询

引擎自动判定查询涉及的 Partition 并路由至归属副本,查询写法与普通表一致,分区裁剪正常生效。
SELECT * FROM events WHERE user_id = 1001;
SELECT count() FROM events WHERE event_date = '2024-01-15';
说明:
选择性复制表的读路由依赖 Analyzer,须保持 enable_analyzer = 1(默认值)。关闭时报错:Selective replication read routing requires enable_analyzer = 1。

本地 JOIN 与 IN

两张选择性复制表的分桶配置完全一致时,同一分桶键值的数据必然落在同一批节点上,JOIN 可在各节点本地完成,无需广播右表;否则自动降级为 GLOBAL 广播执行。
满足本地执行须同时具备以下条件:
序号
条件
1
两侧均配置 SHARD BY,且均非 SHARD RANDOM。
2
两侧 INTO N 取值相同。
3
两侧 SHARD BY 表达式相同。
4
两侧分桶列的类型签名相同。uid UInt64 与 uid String 即使表达式写法一致也不共置,因为哈希结果不同。
5
两侧 replication_factor 相同。
6
两侧均无正在执行的 Partition 迁移。迁移期间保守降级为 GLOBAL。
配置示例:
CREATE TABLE orders (dt Date, uid UInt64, amount Decimal64(2))
ENGINE = ReplicatedMergeTree('/clickhouse/tables/orders', 'r1')
PARTITION BY toYYYYMM(dt)
SHARD BY uid INTO 4
ORDER BY uid
SETTINGS replication_factor = 2;

CREATE TABLE users (dt Date, uid UInt64, name String)
ENGINE = ReplicatedMergeTree('/clickhouse/tables/users', 'r1')
PARTITION BY toYYYYMM(dt)
SHARD BY uid INTO 4
ORDER BY uid
SETTINGS replication_factor = 2;

-- 本地 JOIN
SELECT count() FROM orders o JOIN users u ON o.uid = u.uid;

-- 本地 IN
SELECT sum(amount) FROM orders WHERE uid IN (SELECT uid FROM users);
FULL JOIN 的处理方式不同:未共置的 FULL JOIN 无法按节点并行执行,否则会重复未匹配行,因此降级为在发起节点单点执行。
说明:
上述 6 个条件中任意一项不满足,查询仍会返回正确结果,只是降级为 GLOBAL 广播执行,表现为性能下降而非报错。因此 JOIN 变慢时应先确认执行路径,参见下一节。

确认查询的执行方式

通过 ProfileEvents 判定实际执行路径。
-- 执行目标查询,客户端指定 --query_id 'check_1'
SELECT count() FROM orders o JOIN users u ON o.uid = u.uid;

SYSTEM FLUSH LOGS query_log;

SELECT
ProfileEvents['SelectiveReplicationColocatedJoin'] AS colocated_join,
ProfileEvents['SelectiveReplicationColocatedJoinDegradedToGlobal'] AS degraded_to_global,
ProfileEvents['SelectiveReplicationAutoGlobalRewrite'] AS rewritten_to_global,
ProfileEvents['SelectiveReplicationColocatedIn'] AS colocated_in,
ProfileEvents['SelectiveReplicationFullJoinDegradedToInitiator'] AS full_join_on_initiator
FROM system.query_log
WHERE query_id = 'check_1' AND type = 'QueryFinish'
ORDER BY event_time_microseconds DESC
LIMIT 1;
计数器含义:
计数器
取值大于 0 时的含义
SelectiveReplicationColocatedJoin
JOIN 已本地执行
SelectiveReplicationAutoGlobalRewrite
JOIN 被改写为 GLOBAL,两表未共置
SelectiveReplicationColocatedJoinDegradedToGlobal
两表本可共置,因迁移进行中降级
SelectiveReplicationColocatedIn
IN 子查询已本地执行
SelectiveReplicationColocatedInDegradedToGlobal
IN 未共置,或子查询包含 GROUP BY、DISTINCT、LIMIT、UNION、CTE 等超出可透传范围的结构
结果不符合预期时,按上一节的 6 个条件逐项比对两张表的 SHOW CREATE TABLE 输出。

典型应用场景

下表汇总各场景的关键配置,可先据此定位与您业务匹配的场景。

场景一:大表关联分析

业务背景:订单表与用户表均为亿级规模,需按用户维度关联分析。原生社区方案下,若两表 Sharding Key 不一致,JOIN 会广播右表,网络开销随右表规模线性增长。
方案要点:两表使用完全相同的分桶键、分桶数、replication_factor 与列类型建表。同一 uid 的数据必然落在同一批节点上,JOIN 在各节点本地完成。
CREATE TABLE orders (dt Date, uid UInt64, amount Decimal64(2))
ENGINE = ReplicatedMergeTree
PARTITION BY toYYYYMM(dt)
SHARD BY uid INTO 16
ORDER BY (uid, dt)
SETTINGS replication_factor = 2;

CREATE TABLE users (dt Date, uid UInt64, name String, city String)
ENGINE = ReplicatedMergeTree
PARTITION BY toYYYYMM(dt)
SHARD BY uid INTO 16
ORDER BY (uid, dt)
SETTINGS replication_factor = 2;
关联查询与执行路径确认:
-- 本地 JOIN,无需广播 users 表
SELECT u.city, count() AS order_cnt, sum(o.amount) AS gmv
FROM orders o JOIN users u ON o.uid = u.uid
WHERE o.dt >= '2024-01-01'
GROUP BY u.city;

-- 确认是否真的本地执行(应为 1)
SELECT ProfileEvents['SelectiveReplicationColocatedJoin']
FROM system.query_log WHERE query_id = '<your_query_id>' AND type = 'QueryFinish';
说明:
分桶列的类型必须一致,不只是列名和表达式一致。若一张表的 uid 为 UInt64、另一张为 String,即使 SHARD BY uid INTO 16 写法完全相同,哈希结果也不同,两表不共置。此时查询仍返回正确结果,只是静默降级为 GLOBAL 广播——不报错、只变慢,容易被误判为数据量增长导致的正常劣化。上线后请用 ProfileEvents 确认一次。

场景二:存储成本优化

业务背景:6 节点集群承载明细数据,逻辑数据量 20 TB。使用标准 ReplicatedMergeTree 时每个节点存储全量 20 TB,总占用 120 TB,存储成本高且单机容量成为瓶颈。
方案要点:设置 replication_factor = 2,总占用降至约 40 TB,各节点约 6.7 TB。分桶数按单个业务分区的数据量确定,参见分桶数选择。
CREATE TABLE detail_log
(
event_time DateTime,
event_date Date,
device_id UInt64,
payload String
)
ENGINE = ReplicatedMergeTree
PARTITION BY toYYYYMMDD(event_date)
SHARD BY cityHash64(device_id) INTO 12
ORDER BY (device_id, event_time)
SETTINGS replication_factor = 2;
上线后核对各分区实际数据量,校准后续同类表的分桶数取值:
SELECT partition_id, formatReadableSize(sum(bytes_on_disk)) AS size
FROM system.parts
WHERE database = currentDatabase() AND table = 'detail_log' AND active
GROUP BY partition_id
ORDER BY sum(bytes_on_disk) DESC
LIMIT 20;
说明:
降低 replication_factor 是在用可靠性换存储成本,两者需一并评估。replication_factor = 2 仅能容忍 1 个节点故障;同时故障 2 个节点且恰为某 partition 的全部归属副本时,该分区不可读。核心业务表请评估是否应设置为 3,参见故障容忍度。

场景三:冷热表差异化副本策略

业务背景:同一集群上既有核心交易表(数据量小、可靠性要求高),也有明细日志表(数据量大、可容忍短时不可用)。原生方案下所有表共用一套分片副本布局,无法按表区分。
方案要点:选择性复制表的副本数是表属性,同集群内各表独立设置。核心表设 replication_factor = 3 换取更高容错,容量型表设 2 控制成本。
-- 核心交易表:3 副本,容忍 2 节点故障
CREATE TABLE trade_core (dt Date, order_id UInt64, amount Decimal64(2))
ENGINE = ReplicatedMergeTree
PARTITION BY toYYYYMM(dt)
SHARD BY order_id INTO 6
ORDER BY (order_id, dt)
SETTINGS replication_factor = 3;

-- 明细日志表:2 副本,控制存储成本
CREATE TABLE trade_log (dt Date, order_id UInt64, detail String)
ENGINE = ReplicatedMergeTree
PARTITION BY toYYYYMMDD(dt)
SHARD BY order_id INTO 24
ORDER BY (order_id, dt)
SETTINGS replication_factor = 2;
副本数需要调整时使用专用 DDL,无需重建表:
ALTER TABLE trade_log MODIFY REPLICATION FACTOR 3;
ALTER TABLE trade_log START SELECTIVE REBALANCE SYNC;
说明:
两表的 replication_factor 不同时无法本地 JOIN(本地执行条件之一是两侧 replication_factor 相同)。若这两张表之间存在高频关联查询,需在"差异化副本"与"本地 JOIN"之间取舍:要么统一副本数,要么接受 JOIN 降级为 GLOBAL 执行。

场景四:弹性扩缩容

业务背景:业务数据量季度性增长,需定期扩容。原生方案下扩容要修改集群配置 <remote_servers>,且存量数据保持原分布不自动迁移,新节点只承载新写入数据,负载长期不均。
方案要点:选择性复制表的数据落点由引擎按 Partition 计算,节点加入后触发 Rebalance 即可重分布存量数据。扩缩容通过控制台界面操作,无需修改集群配置。
-- 在控制台扩容后自动触发 Rebalance、同步等待完成,无须人工操作,以下仅列出操作命令
ALTER TABLE detail_log START SELECTIVE REBALANCE SYNC;

-- 观察迁移进度
SELECT state, count() FROM system.selective_migrations
WHERE database = currentDatabase() AND table = 'detail_log'
GROUP BY state;

-- 确认各节点负载已均衡
SELECT role, replica, partition_count
FROM system.selective_replica_load
WHERE database = currentDatabase() AND table = 'detail_log'
ORDER BY partition_count DESC;
迁移期间读写正常可用。迁移速度由 max_concurrent_partition_migrations(默认 16)控制,可用于限速。
说明:
单个物理 partition 过大时,迁移耗时会显著增加,容易触达 migration_timeout_seconds(默认 300 秒)导致 CLONE 阶段回滚。因此分桶数在规划阶段就应把迁移耗时考虑进去,将单个物理 partition 的压缩后数据量控制在 10 GB 以内。命令超时不影响数据一致性,可直接重新下发。

场景五:明细日志高速写入

业务背景:日志采集链路要求写入吞吐优先,数据仅做时间范围检索与聚合统计,不做大表关联,也不需要按 key 去重。
方案要点:使用 SHARD RANDOM。单个 INSERT block 只落入一个分桶,不产生写放大,避免触达 max_partitions_per_insert_block 上限。
CREATE TABLE app_log
(
log_time DateTime,
log_date Date,
level LowCardinality(String),
message String
)
ENGINE = ReplicatedMergeTree
PARTITION BY toYYYYMMDD(log_date)
SHARD RANDOM INTO 12
ORDER BY (log_time, level)
SETTINGS replication_factor = 2;
配合写入参数调优:
SET min_insert_block_size_rows = 1000000,
min_insert_block_size_bytes = 134217728;
说明:
选择 SHARD RANDOM 意味着同时放弃三项能力:写入重试幂等、合并语义引擎、本地 JOIN。其中写入幂等的缺失影响最容易被忽略——采集链路重试时桶号会重新随机生成,无法通过 block 去重识别重复数据,可能产生重复行。若您的采集链路存在重试机制且不能容忍重复,请改用确定性 SHARD BY。

场景六:不适用选择性复制表的场景

业务背景:并非所有表都适合启用选择性复制表。由于选择性复制表默认开启,以下情况需显式关闭,否则会引入预期外的行为。
方案要点:设置 replication_factor = 0,退回标准 ReplicatedMergeTree,所有节点存储全量数据。
-- 显式关闭选择性复制表
CREATE TABLE dim_config (id UInt64, name String, config String)
ENGINE = ReplicatedMergeTree
ORDER BY id
SETTINGS replication_factor = 0;

运维操作

选择性复制表的日常运维围绕 Rebalance 与迁移展开。副本数调整、节点扩缩容、欠副本补齐均通过本节的命令与系统表完成。

运维命令

表级命令:
-- 触发一次 Rebalance,指定 SYNC 时同步等待完成
ALTER TABLE [db.]tbl START SELECTIVE REBALANCE [SYNC];

-- 等待本表所有迁移任务结束
ALTER TABLE [db.]tbl SYNC SELECTIVE MIGRATIONS;

-- 将指定 partition 迁移至指定节点
ALTER TABLE [db.]tbl MIGRATE PARTITION '<partition_id>' TO REPLICA '<replica_name>';

-- 清理超出 replication_factor 的冗余副本,省略 partition_id 表示处理全部分区
ALTER TABLE [db.]tbl TRIM OVER REPLICATED PARTITIONS ['<partition_id>'];

-- 修改指定表的副本数
ALTER TABLE [db.]tbl MODIFY REPLICATION FACTOR N;
节点级命令,作用于该节点上的所有选择性复制表:
-- 标记节点下线并启动后台排空
SYSTEM DECOMMISSION REPLICA '<replica_name>';

-- 取消下线标记
SYSTEM CANCEL DECOMMISSION REPLICA '<replica_name>';
对非选择性复制表执行上述命令时报错:
db.tbl is not a selective-replication ReplicatedMergeTree table (requires replication_factor > 0)
MIGRATE PARTITION 的拒绝场景:
MIGRATE PARTITION requires selective replication enabled
Partition '202401-3' has no assignment
Replica 'r9' is not active
Partition '202401-3' has active migration
START SELECTIVE REBALANCE 的相关错误信息:
Rebalance already in progress on replica r1

SYSTEM START SELECTIVE REBALANCE ...: timed out after N rounds (started=.., skipped=.., failed=..).
In-flight migrations will continue in the background; the ZK lock has been released, you may
safely re-issue the command. See the 'receive_timeout' setting.
说明:
命令超时不影响数据一致性。进行中的迁移会在后台继续执行,分布式锁已释放,可直接重新下发命令。

扩容与缩容

扩容与缩容均无需修改集群配置,通过 控制台界面 操作即可。
操作
流程
扩容
控制台界面执行扩容,系统自动在新节点建表并触发 Rebalance。
缩容
控制台界面执行缩容,系统自动先排空数据再下线节点。
迁移期间读写正常可用,Rebalance 并发度由 max_concurrent_partition_migrations 控制,默认 16,可用于限速。

系统表

system.selective_assignments 记录分区归属关系,是运维排查的主要入口。
字段
类型
说明
database / table
String
库名与表名
partition_id
String
分区 ID
has_partition_value
UInt8
是否携带分区值(v2 格式)。取 0 时分区裁剪保守保留该分区
assigned_replicas
Array(String)
归属副本列表
replicas_by_role
Map(String, Array(String))
按角色分组的归属副本
cloning_replicas
Array(String)
处于迁移克隆中的副本
is_under_replicated
UInt8
是否欠副本,即就绪且活跃的副本数小于 replication_factor
surviving_source_count
UInt64
当前持有该分区的就绪且活跃副本数
zk_read_error
String
非空表示读取 Keeper 失败,该行数据为保守值
system.selective_migrations 记录迁移任务。
字段
类型
说明
database / table
String
库名与表名
migration_id
String
迁移任务 UUID
partition_id
String
迁移的分区
source_replica / target_replica
String
源副本与目标副本
state
String
任务状态,取值为 CLONE、SWITCH、CLEANUP、DONE、FAILED
coordinator
String
协调该任务的副本
snapshot_parts
UInt64
源快照包含的 part 数
created_at
DateTime
创建时间
zk_read_error
String
非空表示读取 Keeper 失败
system.selective_replica_load 记录每个节点的负载。
字段
类型
说明
database / table
String
库名与表名
role
String
节点角色
replica
String
副本名
partition_count
UInt64
该节点在该角色下承载的分区数
role_version
Int32
当前拓扑版本
counts_by_role_version
Int32
统计数据重建时的版本
is_stale
UInt8
两个版本不一致时取 1

监控项

欠副本检查,优先级最高:
SELECT database, table, count() AS under_replicated_partitions
FROM system.selective_assignments
WHERE is_under_replicated
GROUP BY database, table
ORDER BY under_replicated_partitions DESC;
节点负载均衡度:
SELECT role, replica, partition_count
FROM system.selective_replica_load
WHERE database = currentDatabase() AND table = 'events'
ORDER BY partition_count DESC;
分区数据量分布,用于评估分桶键质量与 N 的取值:
SELECT partition, sum(rows) AS rows, formatReadableSize(sum(bytes_on_disk)) AS size
FROM system.parts
WHERE database = currentDatabase() AND table = 'events' AND active
GROUP BY partition
ORDER BY rows DESC
LIMIT 20;
迁移任务状态:
SELECT state, count() FROM system.selective_migrations
WHERE database = currentDatabase() GROUP BY state;
数据完整性告警,建议接入监控系统:
SELECT value FROM system.events WHERE event = 'SelectiveReplicationLostPidsSkipped';
Keeper 读取异常:
SELECT value FROM system.events WHERE event = 'SelectiveReplicationAssignmentZKError';
注意:
SelectiveReplicationLostPidsSkipped 大于 0 表示存在查询在 skip_unavailable_shards = 1 下返回了不完整结果,且未向客户端报错。这是发现该情况的唯一信号,业务开启该参数时必须配置此项告警。同理,SelectiveReplicationAssignmentMissingPids 大于 0 表示归属关系缓存出现一致性异常,须上报处理

可用性与故障处理

选择性复制表的每个 Partition 仅存在于 replication_factor 个节点,可用性模型与标准 ReplicatedMergeTree 不同。业务上线前须确认本节内容。

全部归属副本不可用时的读行为

skip_unavailable_shards 取默认值 0 时直接报错,不返回部分数据。错误码为 ALL_REPLICAS_LOST,编号 415:
Selective replication: 3 partition(s) have no reachable owner (set skip_unavailable_shards=1
to skip). Examples: [202401-3, 202401-5, 202402-1]
设置为 1 时跳过这些分区并返回不完整结果,同时记录 warning 日志并累加 SelectiveReplicationLostPidsSkipped:
Selective replication: skipping 3 partition(s) with no reachable owner
(skip_unavailable_shards=1). Examples: [...]
注意:
这是唯一会导致查询静默返回不完整结果的路径。查询不报错、有结果返回,但数据缺失,下游极易将其当作真实数据使用。不建议长期使用该参数规避故障;确需开启时,必须配置 SelectiveReplicationLostPidsSkipped 计数器告警。

部分归属副本不可用时的读行为

引擎自动选择其他归属副本,无需人工干预。选择顺序:
1. 优先选择同角色的归属副本,跨角色作为备选,累加 SelectiveReplicationCrossRoleRead 与 SelectiveReplicationCrossRoleFallbackOnFailure。
2. selective_replication_route_by_delay 开启时,过滤复制延迟超过 max_replica_delay_for_distributed_queries 的节点。
3. selective_replication_weighted_routing 开启时按实际数据量加权。
4. 处于 :cloning 状态的节点不作为可读副本。

归属副本不可用时的写行为

情况
行为
部分归属副本不可用
写入可用节点,其余节点通过常规复制补齐
全部归属副本不可用
报错,错误信息参见 写入失败的错误信息
配置了 insert_quorum 且可用归属副本不足
返回 TOO_FEW_LIVE_REPLICAS

故障容忍度

以 3 节点集群、replication_factor = 2 为例:
故障节点数
读
写
0
正常
正常
1
正常,每个 Partition 均有另一个归属副本可用
正常,可能临时降为单副本,节点恢复后补齐
2
部分分区不可读。两个故障节点恰为某 Partition 的全部归属副本时,该分区返回 ALL_REPLICAS_LOST
对应分区写入失败
结论:replication_factor = 2 仅可容忍 1 个节点故障。需容忍 2 个节点故障时须设置 replication_factor = 3,且集群节点数不少于 3。
说明:
存储占用与故障容忍度需按表权衡:核心表可提高 replication_factor,容量型表可维持较低取值,两者可共存于同一集群。

故障恢复

节点恢复并重新进入活跃、就绪状态后,自动恢复为归属副本并参与读写。
欠副本不会自动补齐。enable_auto_rebalance 默认关闭,须手动执行 ALTER TABLE ... START SELECTIVE REBALANCE。
长期无法恢复的节点按缩容流程下线。

节点只读的处理

节点的 replication_factor 与 Keeper 记录不一致时,服务启动阶段将该节点置为只读,且不向客户端返回错误。排查方式:
SELECT database, table, replica_name, is_readonly FROM system.replicas WHERE is_readonly;
处理方式:核对该节点建表语句中的 选择性复制表参数,与其他节点保持一致后重启服务。

常见问题

以下汇总建表、读写与运维过程中的高频问题。

没有配置任何 选择性复制表参数,如何判断表是否 选择性复制表?

在弹性版集群上创建的 ReplicatedMergeTree 表引擎默认开启选择性复制能力,replication_factor 默认值为 2 每个 partition 仅存储 2 份数据。若需所有节点存储全量数据,必须显式设置 replication_factor = 0。

replication_factor 和分桶数能修改吗

replication_factor 可以,分桶数不可以。
配置项
可否修改
方式
replication_factor
可以
ALTER TABLE ... MODIFY REPLICATION FACTOR N,之后触发 rebalance
分桶数 INTO N
不可以
只能新建表并迁移数据
分桶键
不可以
只能新建表并迁移数据
请注意 ALTER ... MODIFY SETTING replication_factor 是不支持的路径,其报错文案建议 create a new table,该文案未同步专用 DDL 的支持情况,请以本表为准。

选择性复制表支持 Mutation、TTL、物化视图吗

支持。数据变更(mutation、轻量删除)、生命周期(TTL)、查询(OPTIMIZE FINAL、SELECT ... FINAL)、加速结构(跳数索引、投影、物化视图)四类能力在 选择性复制表上均已验证可用。与标准 ReplicatedMergeTree 的差异在于执行范围:这些操作作用于各 Partition 的归属副本,而非集群全部节点。因此判断执行是否完成时,须按 Partition 逐个确认,不能以单个节点的状态代表全表。

选择性复制表如何做备份

选择性复制表不支持 BACKUP、RESTORE、FETCH PARTITION、FREEZE PARTITION,无法通过原生命令建立表级备份。原因是 选择性复制表的 partition 仅分散在部分节点上,且归属关系记录于 Keeper,上述命令的单机语义无法覆盖这一布局。
可选的替代方案如下,请按业务对恢复目标的要求选择:
方案
做法
适用情况
逻辑导出
INSERT INTO ... SELECT 将数据写入对象存储表引擎或另一集群的表
可接受按分区批量导出,恢复时重新导入
关键表关闭 SR
对必须具备原生备份能力的表设置 replication_factor = 0
表数据量不大,可承受全节点存储全量
提高副本数
将 replication_factor 提升至 3 或以上
以在线冗余替代离线备份,降低单节点故障导致数据丢失的概率
注意:
提高 replication_factor 只能提升在线可靠性,不能替代备份。误删除、误更新等逻辑错误会同步到所有归属副本,此类场景仍需逻辑导出方案兜底。

建表失败怎么排查

错误信息关键字
原因
处理方式
requires N >= 2 / requires N <= 256
分桶数越界
取值范围调整为 [2, 256]
must have an integer type
分桶键非整数类型
使用 cityHash64() 等哈希函数包裹
expects a single expression, got a tuple
分桶键为元组
改为单个表达式
non-integral type ... require every PARTITION BY column
分区键列类型不支持
分区键使用哈希函数包裹,或去掉 SHARD BY
is reserved for the SHARD BY clause
声明了 _shard 列
更换列名
replication_factor mismatch
各节点 SR 参数不一致
统一参数后重建表或重启服务

写入报分区数超限怎么处理

SET max_partitions_per_insert_block = 1000;
调整后仍频繁触发,说明单批数据跨越的业务分区过多或分桶数过大。

查询返回 ALL_REPLICAS_LOST 怎么处理

表示存在分区的全部归属副本不可用。排查方式:
SELECT partition_id, assigned_replicas, surviving_source_count, is_under_replicated
FROM system.selective_assignments
WHERE database = currentDatabase() AND table = 'events' AND is_under_replicated;

SELECT replica_name, is_readonly, absolute_delay FROM system.replicas
WHERE database = currentDatabase() AND table = 'events';
优先恢复故障节点;无法恢复时执行下线流程并触发 rebalance 补齐副本。不建议长期使用 skip_unavailable_shards = 1 规避,该配置会导致查询结果不完整。

JOIN 性能下降怎么排查

先按 确认查询的执行方式 判定实际执行路径,再按 本地 JOIN 与 IN 的 6 个条件逐项比对两表的 SHOW CREATE TABLE 输出。常见原因为分桶数不一致、分桶列类型不一致、replication_factor 不一致、其中一张表存在进行中的迁移。

Rebalance 长时间未完成怎么处理

SELECT state, count() FROM system.selective_migrations
WHERE database = currentDatabase() AND table = 'events' GROUP BY state;

SELECT partition_id, state, source_replica, target_replica, coordinator, created_at
FROM system.selective_migrations
WHERE database = currentDatabase() AND table = 'events' AND state IN ('CLONE', 'SWITCH');
处理要点:
CLONE 阶段超过 migration_timeout_seconds(默认 300 秒)后自动回滚为 FAILED。
协调者失联超过 migration_coordinator_timeout(默认 60 秒)后由其他节点接管。
命令超时不影响数据一致性,可重新下发 ALTER TABLE ... START SELECTIVE REBALANCE。
调大 max_concurrent_partition_migrations 可提升迁移速度,同时增加 Keeper 与网络压力。
单个物理 partition 过大时迁移耗时显著增加。

如何估算存储占用

单表的总存储占用约等于逻辑数据量乘以 replication_factor,与集群节点数无关;各节点承载的数据量约为该值除以节点数。例如 3 节点集群将 replication_factor 从 3 调整为 2 后,总占用约降至原来的三分之二。
由于归属关系按 Rendezvous 哈希计算,各节点的实际数据量存在偏差,分桶数偏小时偏差更明显。

数据倾斜怎么处理

按 监控项 检查节点负载与分区数据量分布。倾斜通常由分桶键基数不足或存在热点值导致,须重建表并更换分桶键。

为什么欠副本没有自动恢复

enable_auto_rebalance 默认为 false。节点恢复后会自动重新参与读写,但历史欠副本的数据补齐需要手动触发:
ALTER TABLE [db.]tbl START SELECTIVE REBALANCE SYNC;