腾讯云数据仓库 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 表引擎以及弹性版之前,请您了解以下限制:
限制项 | 说明 |
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 = KeeperMergeTreeSHARD BY user_id INTO 8ORDER BY (event_time, user_id);
KeeperMergeTree 引擎兼容 MergeTree 语法与功能,在建表时会隐式地将 MergeTree 引擎转换为 KeeperMergeTree:用户无需修改现有 DDL,直接使用
ENGINE = MergeTree 建表即可,引擎会自动替换为 KeeperMergeTree。当然,用户也可以显式指定 ENGINE = KeeperMergeTree,两者等价。-- 以下两种写法等价,最终都会创建为 KeeperMergeTreeCREATE TABLE analytics.events(event_id UInt64,user_id UInt64,event_time DateTime,payload String)ENGINE = MergeTreeSHARD BY user_id INTO 8ORDER BY (event_time, user_id);CREATE TABLE analytics.events(event_id UInt64,user_id UInt64,event_time DateTime,payload String)ENGINE = KeeperMergeTreeSHARD BY user_id INTO 8ORDER 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 过滤:自动裁剪,只命中目标 ShardSELECT event_id, payloadFROM 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 查询:只命中相关 ShardSELECT 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.nameFROM distributed(currentDatabase(), events) AS eJOIN local(currentDatabase(), users) AS u ON e.user_id = u.user_id;-- 右表未使用 local():每个 shard 节点都会全表扫描右表SELECT e.event_id, u.nameFROM distributed(currentDatabase(), events) AS eJOIN 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, payloadFROM 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 建表后不可修改。如确实需要调整,请新建表并迁移数据。因此建表时应按集群最大预计规模设定。弹性版集群适合哪些业务?
数据量大且需要频繁弹性伸缩、存储成本敏感、有明确查询过滤维度的业务最为适合。若数据量小且节点规模固定,标准集群的本地盘方案在延迟上仍有一定优势。