帮你快速理解、总结文档立即下载
文档中心>腾讯云数据仓库 TCHouse-C>开发指南>表引擎>自研存算分离表引擎 KeeperMergeTree

自研存算分离表引擎 KeeperMergeTree

最近更新时间:2026-09-17 21:36:00
我的收藏
腾讯云数据仓库 TCHouse-C 自研了 KeeperMergeTree 表引擎,是弹性版实现存储与计算分离架构的底层基础能力,本文将介绍 KeeperMergeTree 的核心特性、建表方法、数据读写方式、性能优化手段、典型业务场景应用,以及从标准集群迁移的操作要点。

功能优势

在存储与计算分离架构下,数据统一存储在对象存储上,计算节点无状态、支持秒级弹性伸缩。连接任意节点即可查询和写入全量数据,无需手动管理分布式表。通过逻辑分片机制,可进一步获得查询裁剪与缓存亲和性加速。标准版与弹性版的主要差异如下:
对比
标准版
弹性版
架构
存算一体(Shared-Nothing)
存算分离(Shared-Storage)
存储
数据分散在各分片节点本地盘
数据统一存储在对象存储,节点无状态
分片
物理分片,与节点强绑定
逻辑分片,与节点解耦
扩缩容
需搬迁数据,耗时长
秒级弹性伸缩,无需搬迁
DDL / RBAC
需 ON CLUSTER 传播到集群
自动同步所有节点,无需 ON CLUSTER
表引擎
MergeTree 表引擎
KeeperMergeTree
副本
需要显式创建 Replicated 副本表
不需要加 Replicated 前缀,自动创建存算分离表
分布式查询
通过 Distributed 引擎表
默认使用 Parallel Replicas 进行查询(节点级并行),逻辑分片模式可以通过 distributed() 表函数进行查询(逻辑分片级别并行)
写入
写 Distributed 引擎表或本地表
直接写本地表,引擎自动调度
KeeperMergeTree 表引擎在兼容 MergeTree 表引擎的基础上,额外提供以下能力:
特性
说明
高效写入
数据直接写对象存储,提交仅需一次 Keeper 交互,无本地临时 IO。
强一致性
通过 ClickHouse Keeper 保证元数据一致,所有节点最终一致。
秒级弹性
计算节点无状态,启动时从 Keeper 加载元数据即可提供服务。
高可用
至少一个节点存活即可提供服务,数据不因节点下线丢失。
存储高效
Mutation 时未修改的列文件通过引用共享,避免不必要的数据拷贝。
零泄漏
完整的数据生命周期追踪,杜绝对象存储上的孤儿文件。
说明:
当前存算分离表引擎 KeeperMergeTree 处于白名单试用阶段、仅在 弹性版 中支持。若您想要体验该功能,可 联系我们 为您开通白名单。

使用限制

在使用 KeeperMergeTree 表引擎以及弹性版之前,请您了解以下限制:
限制项
说明
shard_num 建表后不可修改
需按集群最大预计规模设定,调整需新建表并迁移数据。
多余节点不参与 Shard 调度
节点数大于 shard_num 时,多余节点仅可通过全量查询 + Parallel Replicas。
长连接存在节点倾斜
VIP 按 TCP 连接做负载均衡,单条长连接始终落在同一节点。
扩缩容导致部分缓存失效
Shard 与节点的映射关系变化后,缓存随后续读写逐步重建。
直接查表不做 Shard 裁剪
直接查询表名返回全量数据,需使用 distributed() 或 local() 才有裁剪收益。

前提条件

已创建 TCHouse-C 弹性版集群(存算分离架构)。
已获取具备目标数据库建表和读写权限的数据库账号。
已通过 VIP 地址或客户端连接到集群。

建表语法

本节主要介绍弹性版集群的建表语法以及 shard_num 与 SHARD BY 这两个关键参数的选择方法。弹性版集群使用 KeeperMergeTree 引擎,由 SHARD BY Shard Key INTO shard_num 将数据按 Shard Key 取模映射到逻辑分片中,具体语法如下:
CREATE TABLE analytics.events
(
event_id UInt64,
user_id UInt64,
event_time DateTime,
payload String
)
ENGINE = KeeperMergeTree
SHARD BY user_id INTO 8
ORDER BY (event_time, user_id);
KeeperMergeTree 引擎兼容 MergeTree 语法与功能,在建表时会隐式地将 MergeTree 引擎转换为 KeeperMergeTree:用户无需修改现有 DDL,直接使用 ENGINE = MergeTree 建表即可,引擎会自动替换为 KeeperMergeTree。当然,用户也可以显式指定 ENGINE = KeeperMergeTree,两者等价。
-- 以下两种写法等价,最终都会创建为 KeeperMergeTree
CREATE TABLE analytics.events
(
event_id UInt64,
user_id UInt64,
event_time DateTime,
payload String
)
ENGINE = MergeTree
SHARD BY user_id INTO 8
ORDER BY (event_time, user_id);

CREATE TABLE analytics.events
(
event_id UInt64,
user_id UInt64,
event_time DateTime,
payload String
)
ENGINE = KeeperMergeTree
SHARD BY user_id INTO 8
ORDER BY (event_time, user_id);
说明:
从标准版 ClickHouse 集群迁移时,DDL 中的 ENGINE = MergeTree 可保持不变,无需手动改为 KeeperMergeTree。

选择 shard_num

SHARD BY ... INTO {shard_num} 中的 shard_num 决定逻辑分片数量,建表后不可修改。不同的节点数与 shard_num 配比会带来不同行为:
节点数 vs shard_num
行为
影响
节点数 < shard_num
单节点承担多个 shard
缓存压力增大,可接受但非理想
节点数 ≈ shard_num
每个 Shard 对应独立节点
理想状态
节点数 > shard_num
多余节点不参与 shard 调度
多余节点可通过全量查询 + Parallel replicas 利用
说明:
建议按业务预期的最终规模设定,而非当前节点数。例如当前 4 节点但未来可能扩到 12 节点,应设为 12。

选择 Shard Key

请选择查询中最常用于过滤或 GROUP BY 的列作为 Shard Key:
业务类型
推荐 Shard Key
效果
用户行为分析
user_id
同一用户的数据集中在一个 Shard
多租户系统
tenant_id
实现租户级隔离查询
日志分析
trace_id
同一条链路的日志集中在一个 Shard

DDL 与权限管理

弹性版集群的所有节点共享元数据,DDL 和 RBAC 操作会自动同步到所有节点,无需 ON CLUSTER。
-- DDL 直接执行,自动全集群生效
CREATE DATABASE analytics;
ALTER TABLE analytics.events ADD COLUMN source String DEFAULT 'unknown';
DROP TABLE analytics.events;

-- RBAC 操作同理
CREATE USER analyst IDENTIFIED WITH sha256_password BY 'password';
CREATE ROLE analyst_role;
GRANT SELECT ON analytics.* TO analyst_role;
GRANT analyst_role TO analyst;
说明:
从标准集群迁移时,移除所有 SQL 中的 ON CLUSTER 子句即可。

写入数据

直接写入本地表

写入时直接 INSERT 到目标表,引擎会自动按 Shard 归属将数据调度到正确节点,无需借助分布式表。
INSERT INTO analytics.events VALUES
(1, 100, '2025-01-01 00:00:00', 'click'),
(2, 200, '2025-01-01 00:00:01', 'view'),
(3, 100, '2025-01-01 00:00:02', 'purchase');

连接方式

请通过腾讯云提供的 VIP(负载均衡) 连接集群,请求会自动分发到健康节点。
方式
说明
推荐度
连接池
多条连接经 VIP 自然打散到不同节点
推荐
直连节点
查询 system.clusters 获取 default cluster 的节点列表,手动分配写入目标
特定场景
说明:
VIP 按 TCP 连接做负载均衡,单条长连接会始终落在同一节点,导致写入倾斜。请优先使用连接池方式,让多条连接自然打散。

查询数据

弹性版集群提供三种查询模式,核心区别在于访问范围与是否做 Shard 裁剪:
访问方式
执行范围
Shard 裁剪
返回数据
适用场景
distributed(db, table)
跨节点分布式并行
是
所有 Shard 的完整结果
替代 Distributed 表,生产主力
local(db, table)
仅当前节点
是
当前节点 Shard 的数据
调试、单 Shard 查询
直接查 table
仅当前节点
否
全量数据
全表扫描,配合 Parallel Replicas
三者的命名逻辑如下:
distributed():对标 ClickHouse 的 Distributed 引擎,语义完全一致——将请求扇出到各 Shard 节点并行执行并汇总完整结果。从 Distributed 表迁移过来零学习成本。
local():本地 Shard 视图,仅在当前节点执行,只返回归属于本节点 Shard 的数据,不发起跨节点请求。
直接查表:不做任何 Shard 裁剪,返回全量数据(因为每个节点都持有全量元数据)。

分布式查询:distributed()

-- 完整结果:跨所有 Shard 并行
SELECT count(), sum(user_id)
FROM distributed(currentDatabase(), events);

-- 按 Shard Key 过滤:自动裁剪,只命中目标 Shard
SELECT event_id, payload
FROM distributed(currentDatabase(), events)
WHERE user_id = 10086;

-- 按 Shard Key GROUP BY:每个分组完整落在单个 Shard,无需跨节点聚合
SELECT user_id, count()
FROM distributed(currentDatabase(), events)
GROUP BY user_id;

-- IN 查询:只命中相关 Shard
SELECT count()
FROM distributed(currentDatabase(), events)
WHERE user_id IN (1, 2, 3, 100);

JOIN 查询

当左表是 distributed() 时,右表的行为取决于是否使用 local():
右表使用 local():右表会按照左表 distributed() 底表的 shard 归属来分发查询——每个 shard 节点只查询自己 shard 的右表数据,与左表对齐,避免全表扫描。
右表未使用 local():右表查询全量数据,每个 shard 节点都会对右表执行一次全表扫描。
-- 右表显式使用 local():按左表 shard 分发,避免全表扫描
SELECT e.event_id, u.name
FROM distributed(currentDatabase(), events) AS e
JOIN local(currentDatabase(), users) AS u ON e.user_id = u.user_id;

-- 右表未使用 local():每个 shard 节点都会全表扫描右表
SELECT e.event_id, u.name
FROM distributed(currentDatabase(), events) AS e
JOIN users AS u ON e.user_id = u.user_id;
说明:
当右表也建有 SHARD BY 且 Shard Key 与 JOIN 条件一致时,右表使用 local() 可显著减少扫描量;否则右表只能全表查询。

本地 Shard 查询:local()

local() 支持可选的 shard_num 参数,支持单值或数组:
省略 shard_num:自动推算当前节点所属的 Shard,仅返回该 Shard 的数据。
指定单个 shard_num:查询指定的逻辑分片,不限于本节点负责的 Shard。
指定 shard_num 数组:同时查询多个逻辑分片,返回这些 Shard 的合并结果。
-- 当前节点 Shard 的数据量
SELECT count()
FROM local(currentDatabase(), events);

-- 本地 Shard 内按 Shard Key 查询(命中则返回,未命中返回空)
SELECT event_id, payload
FROM local(currentDatabase(), events)
WHERE user_id = 10086;

-- 指定 shard_num,查询第 4 号逻辑分片的数据
SELECT count()
FROM local(currentDatabase(), events, 4);

-- 指定数组,查询第 1、3、5 号逻辑分片的数据
SELECT count()
FROM local(currentDatabase(), events, [1, 3, 5]);

全量查询 + Parallel Replicas

-- 直接查本地表,任意节点返回全量数据
SELECT count() FROM events;

-- 配合 parallel replicas 在多节点并行加速全表扫描
SET enable_parallel_replicas = 1;
SELECT count(), sum(user_id) FROM events;

表函数签名

distributed(database, table)
local(database, table [, shard_num])
database:数据库名,支持 currentDatabase() 引用当前库;
table:表名;
shard_num:local() 的可选参数,指定要查询的逻辑分片编号(从 1 开始)。支持单个整数或数组,省略时自动推算当前节点所属的 Shard。
说明:
两个表函数默认开启 Shard 裁剪,无需额外设置。当 WHERE 条件包含 Shard Key 的等值或 IN 过滤时,引擎会自动跳过不相关的 Shard。

技术原理

本节介绍逻辑分片和缓存亲和性的技术原理,供用户更深入地理解系统底层的实现机制。

逻辑分片 Shard

每个 Data Part 在写入时计算 Shard Value,按以下规则归属到逻辑分片:
逻辑 Shard 编号 = hash(shard_key_value) % shard_num
例如表定义为 SHARD BY user_id INTO 12 时,user_id = 100 的数据归属到 100 % 12 = 4 号 shard。集群启动时,逻辑 Shard 按轮询方式分配到物理节点:
节点数 / Shard 数
分配结果
4 节点 / 12 Shard
每节点负责 3 个 Shard
12 节点 / 12 Shard
每节点负责 1 个 Shard(理想状态)
16 节点 / 12 Shard
12 个节点各负责 1 个 Shard,4 个节点不参与 Shard 调度
当 distributed() 或 local() 查询的 WHERE 条件包含 Shard Key 的等值或 IN 过滤时,引擎会计算目标 Shard 编号,只将请求发送到对应节点,跳过无关 Shard。

分片缓存亲和性 Shard Affinity

存算分离架构的数据存储在对象存储上,通过本地 Filesystem Cache 缓存热数据来加速读取。如果写入和后台任务(Merge、Mutate)的调度不按 Shard 归属对齐,同一份数据会在不同节点间漂移,导致缓存反复失效。因此 Shard Affinity 分片缓存亲和性让数据写入、后台任务、查询都在同一节点上形成稳定的亲和关系,从而最大化缓存命中率。
Shard Affinity 配置 enable_shard_affinity 集群内默认开启,无需逐表指定。
<merge_tree>
<enable_shard_affinity>1</enable_shard_affinity>
</merge_tree>

缓存失效场景

场景
影响
说明
不使用 distributed() / local() 查询
无亲和性
直接查本地表走全量读取,缓存按默认调度
未开启 enable_shard_affinity
无亲和性
写入和后台任务不按 shard 对齐
集群扩缩容
部分缓存失效
Shard 与节点的映射变化,缓存随后续读写逐步重建(预期行为)

实践教程

以下建议可帮助您在生产环境中获得更优的性能,并提升可维护性。

按最终规模设定 shard_num

shard_num 建表后不可修改,请按业务预期的最大节点数设定。设置偏小会导致后续扩容时多余节点无法参与 Shard 调度,设置过大则单节点承担的 Shard 过多、缓存压力增加。

统一使用 distributed() 查询

生产业务查询应统一使用 distributed() 表函数,以同时获得 Shard 裁剪与并行加速收益。local() 建议仅用于调试 Shard 归属或构建节点级聚合 Pipeline。

排序键以 Shard Key 开头

建议将 Shard Key 放在 ORDER BY 的第一位(如 ORDER BY (user_id, event_time)),使单 Shard 内数据物理聚集,进一步减少扫描量。

使用连接池避免写入倾斜

VIP 按 TCP 连接做负载均衡,单条长连接会始终落在同一节点。请使用连接池让多条连接自然打散;若需精确控制,可查询 system.clusters 获取节点列表后手动分配。

关注热点 Shard

若 Shard Key 的数据分布严重不均(如个别超大租户),对应 Shard 会成为热点。此时可考虑将超大主体单独拆表,或使用复合表达式作为 Shard Key 打散数据。

常见问题

直接查本地表和使用 distributed() 有什么区别?

直接查本地表返回全量数据,不做 Shard 裁剪,等价于全表扫描。distributed() 将请求路由到各 Shard 对应节点并行执行,每个节点只读自己 Shard 的数据,既有裁剪收益又有并行加速。

local() 和直接查表有什么区别?

local() 会做 Shard 裁剪,只返回归属于当前节点 Shard 的数据;直接查表不裁剪,返回全量数据。两者都在本地执行、不跨节点。

什么时候应该用 local() 而不是 distributed()?

local() 只在本节点执行,适合调试 Shard 归属、构建节点级别的聚合管道,或确定只需要本 Shard 数据的场景。大多数业务查询应使用 distributed()。

扩容后需要做什么操作?

无需任何操作。新节点会自动加载全量元数据并开始服务。Shard 与节点的映射会重新计算,部分缓存将逐步失效并重建,属于预期行为。

可以不使用 SHARD BY 建表吗?

可以。不指定 SHARD BY 时表没有逻辑分片,所有节点查询均返回全量数据,可通过 Parallel Replicas 实现多节点并行加速。

shard_num 设置错误怎么办?

shard_num 建表后不可修改。如确实需要调整,请新建表并迁移数据。因此建表时应按集群最大预计规模设定。

弹性版集群适合哪些业务?

数据量大且需要频繁弹性伸缩、存储成本敏感、有明确查询过滤维度的业务最为适合。若数据量小且节点规模固定,标准集群的本地盘方案在延迟上仍有一定优势。