分布式表与集群
单机 ClickHouse 已经能轻松吃下几十亿行,但终有上限:磁盘不够、CPU 打满、需要容灾。这一章讲如何把多台机器组成一个集群,用分片扩展容量、用副本保障高可用,并通过 Distributed 引擎像查一张表一样查整个集群。
1. 两个正交的概念:分片与副本
很多人把"分片"和"副本"混为一谈,其实它们解决完全不同的问题:
- 分片(Shard):把数据水平切分到不同机器上,目的是扩展容量和吞吐。分片之间是"各存一部分",互不包含。
- 副本(Replica):同一份数据的多份拷贝,放在不同机器上,目的是高可用和读扩展。副本之间"内容相同"。
一个集群 = 多个分片,每个分片 = 1 个或多个副本
shard 1 shard 2
┌─────────────┐ ┌─────────────┐
│ replica 1a │ │ replica 2a │
│ replica 1b │ │ replica 2b │
└─────────────┘ └─────────────┘
存 user_id 0~4999 存 user_id 5000~9999
的数据(分片1) 的数据(分片2)
分片解决"装不下/算不动" 副本解决"挂了怎么办"你可以有"2 分片 × 1 副本"(纯扩容无容灾),也可以有"1 分片 × 3 副本"(纯容灾无扩容),更常见的是"N 分片 × 2 副本"。ClickHouse 的副本依赖外部协调服务(ZooKeeper / ClickHouse Keeper,详见 18 章),而分片的路由由 Distributed 引擎完成。
和 MySQL 的"主从复制"对比:MySQL 主从通常是为了读写分离或备份,数据全量相同;ClickHouse 的分片是真的把数据切开,单台机器只持有部分数据,所以跨分片查询必须由协调节点(通常是发起查询的那个节点)汇总。
2. 定义集群:config.d/ 里的 remote_servers
集群拓扑写在配置文件里。ClickHouse 24.x 推荐把所有自定义配置放到 config.d/ 目录(覆盖 config.xml),不要直接改 config.xml。
<!-- /etc/clickhouse-server/config.d/cluster.xml -->
<clickhouse>
<remote_servers>
<my_cluster>
<!-- 第一个分片,含两个副本 -->
<shard>
<replica>
<host>ch-node1</host>
<port>9000</port>
</replica>
<replica>
<host>ch-node2</host>
<port>9000</port>
</replica>
</shard>
<!-- 第二个分片,含两个副本 -->
<shard>
<replica>
<host>ch-node3</host>
<port>9000</port>
</replica>
<replica>
<host>ch-node4</host>
<port>9000</port>
</replica>
</shard>
</my_cluster>
</remote_servers>
<!-- 各节点必须能解析彼此的主机名,并确保 9000(内部)、8123(HTTP)、9009(副本) 端口互通 -->
</clickhouse>配置好后,用 system.clusters 验证:
SELECT cluster, shard_num, replica_num, host_name
FROM system.clusters
WHERE cluster = 'my_cluster';┌─cluster────┬─shard_num─┬─replica_num─┬─host_name─┐
│ my_cluster │ 1 │ 1 │ ch-node1 │
│ my_cluster │ 1 │ 2 │ ch-node2 │
│ my_cluster │ 2 │ 1 │ ch-node3 │
│ my_cluster │ 2 │ 2 │ ch-node4 │
└────────────┴───────────┴─────────────┴───────────┘3. Distributed 引擎:一张"逻辑表"
关键认知:Distributed 表本身不存任何数据,它只是一个路由层。背后每个分片的每个副本上,都要有一张真正存储数据的本地表(通常是 MergeTree 或 ReplicatedMergeTree)。
标准套路是"本地表 + 分布式表"两层:
-- 1) 在每一台节点的本地,建真正的存储表(本地表通常用 _local 后缀)
CREATE TABLE events_local ON CLUSTER my_cluster
(
event_time DateTime,
user_id UInt64,
event_type LowCardinality(String),
page String,
country LowCardinality(String),
device LowCardinality(String),
duration_ms UInt32
)
ENGINE = MergeTree
PARTITION BY toYYYYMM(event_time)
ORDER BY (country, event_type, event_time);
-- 2) 建一个 Distributed 逻辑表,指向集群 + 本地表名
CREATE TABLE events ON CLUSTER my_cluster
AS events_local
ENGINE = Distributed(my_cluster, default, events_local, rand());Distributed 的四个参数:
- 集群名(
my_cluster) - 本地表所在数据库(
default) - 本地表名(
events_local) - 分片键(最后一个参数,如
rand()、user_id、cityHash64(user_id))——决定每行落到哪个分片
CREATE TABLE xxx ON CLUSTER my_cluster 会把这条建表语句广播到集群所有节点,不用你手动在每台机器上敲一遍。注意:副本之间的数据复制不靠 ON CLUSTER,而靠 ReplicatedMergeTree(18 章);ON CLUSTER 只负责"元数据/DDL 的批量下发"。
4. 分布式写入流程
当你 INSERT INTO events(分布式表)时,发生了什么?
客户端 → 连接任意节点(假设 node1)
│
│ 收到 INSERT
▼
node1 按分片键计算每行目标分片
│
┌──────────┴───────────┐
▼ ▼
shard1 的行 shard2 的行
│ │
直接本地写入 通过网络发往 node3/node4 写入
events_local events_local- 如果写入命中本节点所在的分片(比如 node1 属于 shard1),ClickHouse 直接本地落盘,不发网络。
- 如果命中远程分片,会通过 9000 端口把数据转发过去再落盘。
-- 推荐:客户端直接批量写入,由 Distributed 自动按分片键路由
INSERT INTO events
SELECT
now() - number,
number % 100000,
['click','view','purchase'][number % 3 + 1],
concat('/p/', toString(number % 500)),
['CN','US','JP'][number % 3 + 1],
['ios','android','web'][number % 3 + 1],
number % 1000
FROM numbers(1000000);ClickHouse 的写入单位是"批"。每条 INSERT 至少生成一个新的 data part,单条小批量会导致 part 数量爆炸、后台 merge 吃紧,触发 Too many parts 报错。无论是直写分布式表还是本地表,都要攒够一批(几万到几十万行)再写一次。详见 10 章。
生产上更稳妥的做法是:客户端按分片键自己计算目标分片,直接连到对应节点写入 events_local。这样写入路径最短、不会经过转发节点成为瓶颈。Distributed 表则只用于"读"。
5. 分布式查询流程
当你 SELECT ... FROM events(分布式表)时,流程是:
客户端 → node1(协调节点)
│
│ 1. 把查询改写,下发到各分片的"一个副本"
▼
shard1 选 replica 1a shard2 选 replica 2a
│ │
│ 2. 各分片本地执行 │
▼ ▼
返回本分片的部分结果 返回本分片的部分结果
│ │
▼ 3. node1 汇总合并 ◀──┘
│
▼
返回给客户端几个要点:
- 协调节点默认只向每个分片的某一个副本发查询(由
load_balancing策略决定选哪个,如random/nearest_hostname/in_order),不会把查询发给分片的所有副本。 - 聚合类查询(如
COUNT、SUM)会"下推"到各分片先做部分聚合,协调节点只合并中间结果,网络传输量很小。 ORDER BY ... LIMIT这类查询,各分片先排好取 top-N,协调节点再合并排序取全局 top-N。
-- 查询会被自动下推到各分片,协调节点只合并
SELECT country, count() AS pv
FROM events
WHERE event_time >= '2024-01-01'
GROUP BY country
ORDER BY pv DESC;6. 分片键怎么选(sharding key)
分片键决定数据如何分布,是分布式设计里最关键的选择。
分片键选得好:
- 数据在各分片间均匀(避免"数据倾斜",某分片撑爆)
- 常用查询能"定位到少数分片",避免全分片广播
分片键选得差:
- 某分片特别大(如按 country 分片,CN 行数远超 JP)
- 几乎每个查询都要扫所有分片常见选择:
| 分片键 | 适用 | 注意 |
|---|---|---|
rand() | 纯扩容、无需按列定位 | 数据完全随机,没法只查某个用户所在分片 |
cityHash64(user_id) | 按用户聚合分析 | 同一用户落到同一分片,关联查询不跨分片 |
toYYYYMM(event_time) | 按时间分布 | 容易倾斜(近期分片很热),通常只用于分区分片一致的场景 |
分片键写在 Distributed 表定义里。如果建错了,只能新建一张分布式表、重新导数据。前期务必想清楚"你的查询大多按哪个维度聚合",让分片键与之对齐。例如用户行为分析大多按 user_id 聚合,就用 cityHash64(user_id)。
7. 跨分片 IN / JOIN:一定要用 GLOBAL
这是分布式查询最容易踩的坑,也是性能杀手。
问题:普通 IN (SELECT ... FROM 分布式表) 或 JOIN 分布式表,在分布式环境下,"子查询/右表"会在每个分片上各自重新执行一次,而且每个分片执行时又会去查所有分片——形成 N×N 的查询风暴。
-- 危险写法:右表 events 是分布式表,每个分片都会对 events 再发起一次全集群查询
SELECT *
FROM events AS a
WHERE a.user_id IN (SELECT user_id FROM events WHERE event_type = 'purchase');
-- 4 个分片 → 每个分片发起 4 次子查询 → 16 次查询,且随分片数平方级恶化解法:用 GLOBAL IN / GLOBAL JOIN。语义是"先由协调节点把子查询结果(或右表)算出来,作为一个临时表广播给每个分片,各分片直接拿这份已算好的数据做本地 IN/JOIN"。
-- 正确写法:GLOBAL 把子查询结果在协调节点算一次,再下发
SELECT *
FROM events AS a
WHERE a.user_id GLOBAL IN (SELECT user_id FROM events WHERE event_type = 'purchase');普通 IN: 分片1查全集群、分片2查全集群、分片3查全集群 ... → N×N
GLOBAL IN:协调节点算一次结果集 → 广播给分片1/2/3 ... → N+1GLOBAL 把子查询结果在协调节点物化后广播。如果子查询结果很大(比如几千万行),协调节点内存会吃紧。此时应改用 15 章讲的"字典"或"预聚合宽表"来替代大表 JOIN。小表/中表用 GLOBAL 没问题。
8. 小结
- 分片扩容、副本容灾,两者正交;一个集群 = 多个分片,每分片 = 1+ 副本。
- 集群拓扑写在
config.d/的remote_servers;system.clusters可验证。 Distributed表只是路由层,不存数据;背后每张本地表才是真存储,标准写法是"本地表 + 分布式表"。- 写入按分片键路由(命中本节点直写,否则转发);查询由协调节点汇总,聚合下推。
- 分片键选
cityHash64(user_id)这类与查询维度对齐的键,避免倾斜。 - 跨分片
IN/JOIN务必用GLOBAL,否则查询量随分片数平方级恶化。
- 用 2 个 Docker 容器(或
clickhouse-server多实例)搭一个2 分片 × 1 副本的my_cluster,用config.d/cluster.xml声明,并SELECT * FROM system.clusters验证。 - 建
events_local(MergeTree)和events(Distributed(my_cluster, default, events_local, rand())),插入 100 万行,分别SELECT count()和按country聚合,观察结果是否等于各分片之和。 - 构造一个跨分片的
IN (SELECT ...)与GLOBAL IN (SELECT ...),用system.query_log的read_rows字段对比两者的扫描量差异(提示:开两个会话分别跑,再用SELECT query, read_rows FROM system.query_log ORDER BY event_time DESC LIMIT 4)。