Learn
ClickHouse/17-distributed

分布式表与集群

单机 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))——决定每行落到哪个分片
ℹ️ON CLUSTER 让 DDL 在全部节点执行

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);
⚠️不要在单条 INSERT 里「逐行」写

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+1
💡GLOBAL 的代价是「右表/子查询」结果要放进内存

GLOBAL 把子查询结果在协调节点物化后广播。如果子查询结果很大(比如几千万行),协调节点内存会吃紧。此时应改用 15 章讲的"字典"或"预聚合宽表"来替代大表 JOIN。小表/中表用 GLOBAL 没问题。

8. 小结

  • 分片扩容、副本容灾,两者正交;一个集群 = 多个分片,每分片 = 1+ 副本。
  • 集群拓扑写在 config.d/ 的 remote_servers;system.clusters 可验证。
  • Distributed 表只是路由层,不存数据;背后每张本地表才是真存储,标准写法是"本地表 + 分布式表"。
  • 写入按分片键路由(命中本节点直写,否则转发);查询由协调节点汇总,聚合下推。
  • 分片键选 cityHash64(user_id) 这类与查询维度对齐的键,避免倾斜。
  • 跨分片 IN/JOIN 务必用 GLOBAL,否则查询量随分片数平方级恶化。
🎯动手练习
  1. 用 2 个 Docker 容器(或 clickhouse-server 多实例)搭一个 2 分片 × 1 副本的 my_cluster,用 config.d/cluster.xml 声明,并 SELECT * FROM system.clusters 验证。
  2. 建 events_local(MergeTree)和 events(Distributed(my_cluster, default, events_local, rand())),插入 100 万行,分别 SELECT count() 和按 country 聚合,观察结果是否等于各分片之和。
  3. 构造一个跨分片的 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)。