Skip to content

第 13 章 副本与分布式(灵魂章)

学习目标:彻底分清「副本(Replica)= 镜像备份」与「分片(Shard)= 横向切割」这两个最容易被新人混淆的概念;能在白板上画出 3 分片 × 2 副本集群下「INSERT 写入 + SELECT 查询」的完整数据流;能用 ReplicatedMergeTree + Distributed 这套黄金搭档自己搭出一个高可用 OLAP 集群;能讲清楚 ClickHouse Keeper 替代 ZooKeeper 的原因,能用 system.replicas / system.replication_queue 排查同步异常,并知道为什么 ClickHouse「没有自动 rebalance」是个坑。


0. 开篇:4 张图看懂部署形态

ClickHouse 的部署形态从「单机」到「多分片多副本」一共 4 档。在动手写引擎之前,先用一张图把它们的差异钉死:

┌──────────────────────────────────────────────────────────────────────┐
│  ① 单机                                                                │
│  ─────                                                                 │
│        ┌────────────┐                                                  │
│        │ CK 节点 N1 │   全量数据:1 亿行                                │
│        └────────────┘                                                  │
│  适合:单业务线 / POC / 数据量 < 10 亿                                  │
│  缺点:宕机即不可用;写入吞吐受单机限制                                  │
├──────────────────────────────────────────────────────────────────────┤
│  ② 单分片 · 多副本(高可用,不扩容量)                                   │
│  ─────                                                                 │
│        ┌────────────┐    ┌────────────┐                                │
│        │ N1 副本 r1 │ ⇄ │ N2 副本 r2 │   两边数据一模一样:1 亿行       │
│        └────────────┘    └────────────┘                                │
│  适合:数据量不大但要求 7×24                                            │
│  优点:任何一台挂了,另一台立刻能扛                                      │
│  缺点:总数据量没增长(两台机器存的是同一份)                            │
├──────────────────────────────────────────────────────────────────────┤
│  ③ 多分片 · 单副本(扩容量,不高可用)                                   │
│  ─────                                                                 │
│        ┌────────────┐  ┌────────────┐  ┌────────────┐                  │
│        │ N1 分片 s1 │  │ N2 分片 s2 │  │ N3 分片 s3 │                   │
│        │ 1/3 数据   │  │ 1/3 数据   │  │ 1/3 数据   │                   │
│        └────────────┘  └────────────┘  └────────────┘                  │
│  适合:纯实验环境 / 离线导数                                            │
│  优点:总容量 ×3,查询并行 ×3                                           │
│  缺点:任何一台挂了 → 1/3 数据不可读                                    │
├──────────────────────────────────────────────────────────────────────┤
│  ④ 多分片 · 多副本(生产标配)⭐                                          │
│  ─────                                                                 │
│   shard1 ┌────────────┐ ┌────────────┐                                  │
│          │ N1 r1      │ │ N2 r2      │   1/3 数据 × 2 份                │
│          └────────────┘ └────────────┘                                  │
│   shard2 ┌────────────┐ ┌────────────┐                                  │
│          │ N3 r1      │ │ N4 r2      │   1/3 数据 × 2 份                │
│          └────────────┘ └────────────┘                                  │
│   shard3 ┌────────────┐ ┌────────────┐                                  │
│          │ N5 r1      │ │ N6 r2      │   1/3 数据 × 2 份                │
│          └────────────┘ └────────────┘                                  │
│  适合:生产环境标准答案                                                  │
│  特点:扩容量 + 高可用,6 节点存 2 倍数据,丢任意一台仍能跑              │
└──────────────────────────────────────────────────────────────────────┘

把这张图刻进脑子里。后面所有 SQL、所有配置、所有动画都是在围绕这张图展开。

📌 金句先记住:「副本是镜像,分片是切割。 副本解决『一台挂了怎么办』,分片解决『一台装不下怎么办』。两者没有谁替代谁,生产里通常一起用。」


1. 副本与分片:先彻底分清楚

1.1 用「图书馆」打比方

比喻一:副本(Replica) = 同一本书的两本印刷件
─────────────────────────────────────────────
图书馆给《三体》买了 2 本,分别放在 1 楼和 3 楼。
内容一字不差,互为备份。
任何一本丢了,另一本完整顶上。

比喻二:分片(Shard) = 一本书拆成多卷
─────────────────────────────────────────────
《史记》100 卷太厚装不下一个柜子,拆成 3 部分:
  - 卷 1-33 放 A 柜
  - 卷 34-66 放 B 柜
  - 卷 67-100 放 C 柜
任何一卷丢了,整本《史记》就残了。

比喻三:分片 + 副本 = 拆成 3 部分,每部分各印 2 本
─────────────────────────────────────────────
1-33 卷各 2 本(A1、A2 柜)
34-66 卷各 2 本(B1、B2 柜)
67-100 卷各 2 本(C1、C2 柜)
丢任意一柜都能补上。

新人最常犯的错有两个:

  1. 把副本当扩容:以为加一个副本,集群存的数据就翻倍了 —— 错!副本只是镜像,总容量没变
  2. 把分片当备份:以为加分片就高可用了 —— 错!分片只是切割,任意一片挂了,那一片的数据就读不到了

1.2 概念对照表

维度副本 Replica分片 Shard
解决的问题高可用、读扩展容量扩展、写扩展
数据关系完全相同的镜像互不相交的切片
总数据量不变(×1)变大(×N)
节点挂了的影响其他副本顶上,业务无感该分片数据不可读,需重启或恢复
实现引擎ReplicatedMergeTreeDistributed(路由) + 各节点本地 MergeTree
协调依赖ZooKeeper / ClickHouse Keeper仅依赖配置文件中的 remote_servers
MySQL 类比主从复制 / Group Replication分库分表 / Vitess / TiDB shard

2. ReplicatedMergeTree:副本同步的核心引擎

2.1 一句话定义

ReplicatedMergeTreeMergeTree 引擎的「带同步功能的孪生兄弟」:每次本地有 Part 落盘 / Merge / Mutation,都会在 ZooKeeper(或 Keeper)里登记一笔,其它副本看到登记后主动拉取并 fetch 对应的 Part 文件。

2.2 「双胞胎照镜子」类比

左副本 N1                           右副本 N2
┌────────────┐                    ┌────────────┐
│ INSERT 来了 │                    │            │
│  ↓         │                    │            │
│ 写本地 Part │                    │            │
│  ↓         │                    │            │
│ 把元信息    │ ──→  ZooKeeper  ──→ │ 收到通知   │
│ 登记到 ZK  │      /clickhouse    │ 主动 fetch  │
│            │      /tables/...    │ 复制 Part  │
│            │                    │  ↓         │
│            │                    │ 写本地     │
│            │                    │ 与 N1 一致 │
└────────────┘                    └────────────┘

ZooKeeper / Keeper 在这里扮演的角色是 「公告板」

  • 谁写了什么 Part(log 日志)
  • 谁正在 Merge 哪些 Part(assignments)
  • 谁还没拿到这条记录(queue 队列)

副本之间不直接握手,全靠这块「公告板」协调。这是 ClickHouse 选择「日志驱动 + 拉取式」复制的关键设计。

2.3 建表 DDL 与宏(Macros)

sql
CREATE TABLE learn_ck.events_local ON CLUSTER ck_cluster
(
    event_time  DateTime,
    user_id     UInt64,
    event_type  LowCardinality(String),
    properties  String
)
ENGINE = ReplicatedMergeTree(
    '/clickhouse/tables/{shard}/events_local',  -- ZK 路径,按 shard 区分
    '{replica}'                                  -- 当前副本名
)
PARTITION BY toYYYYMM(event_time)
ORDER BY (event_type, user_id, event_time);

注意两个关键宏

含义谁来填
{shard}当前节点属于哪个分片(如 010203config.xml<macros>
{replica}当前节点是该分片的哪个副本(如 replica_areplica_bconfig.xml<macros>

config.xml 里这样写:

xml
<!-- 节点 N1:shard 1 的 replica a -->
<macros>
    <shard>01</shard>
    <replica>replica_a</replica>
</macros>

这样 ON CLUSTER 一发,每个节点把自己的 {shard} {replica} 替换进 ZK 路径,就自然形成了:

/clickhouse/tables/01/events_local   → 由 N1(replica_a) + N2(replica_b) 共同维护
/clickhouse/tables/02/events_local   → 由 N3(replica_a) + N4(replica_b) 共同维护
/clickhouse/tables/03/events_local   → 由 N5(replica_a) + N6(replica_b) 共同维护

📌 首次术语解释 · ON CLUSTER:CK 提供的 DDL 分发语法。在任何一个节点执行 CREATE TABLE ... ON CLUSTER ck_cluster (...),CK 会把这条 DDL 推送到 ck_cluster 集群中的所有节点上各自执行一次。底层依赖 ZK 的 DDL 队列(/clickhouse/task_queue/ddl)。

2.4 副本同步流程图

┌──────────────────────────────────────────────────────────────────┐
│                    一次 INSERT 的完整副本同步                      │
└──────────────────────────────────────────────────────────────────┘

时间 ─────────────────────────────────────────────────────────────→

[N1 副本_a]  收到 INSERT                                              

              ├─① 排序 + 列拆分 + 压缩 → 落本地 Part `20260417_5_5_0`  

              ├─② 在 ZK 里 createSequential(/log/log-)               
              │     payload = {"type":"get_part","name":"20260417_5_5_0"}

              └─③ 客户端收到 OK,写入完成                              

[ZooKeeper]  /clickhouse/tables/01/events_local/log/log-0000000123 
              内容:{"type":"get_part","name":"20260417_5_5_0",       
                     "source_replica":"replica_a", "block_id":"..."}  

[N2 副本_b]  ⏰ 后台 ReplicatedMergeTreeQueueTask 周期性轮询          

              ├─① 发现 log/log-0000000123 → 复制到自己的 queue        

              ├─② 从队首取出任务,决定从 source_replica=replica_a 拉取

              ├─③ 走 HTTP 9009(interserver port)下载完整 Part 文件夹
              │     包含 .bin / .mrk2 / primary.idx / checksums.txt   

              ├─④ 校验 checksum → 写本地 → 提交                       

              └─⑤ 在 ZK 标记任务完成(queue 出队)                     
                                                                      
最终:N1 与 N2 都拥有 `20260417_5_5_0`,互为镜像。

注意几个关键点:

  1. 拉取式而非推送式:N2 主动从 N1 拉,而不是 N1 主动推过来。这样写入端的延迟不受副本数影响。
  2. 登记的是元信息,不是数据:ZK 里只放「需要做什么」的 JSON,几百字节;真正的几百 MB 的列文件走 N2→N1 之间的 HTTP(默认 9009 端口)传输。
  3. Merge 也要走 ZK:每次后台 Merge 选中哪些 Part 合成一个新 Part,需要 leader 副本(每个分片选一个)在 ZK 里登记 MERGE_PARTS 任务,其他副本同样从 queue 拉取。
  4. 写入幂等:同一个 INSERT 如果客户端重试,CK 通过 block_id(数据块的哈希)去重,不会写重复数据。

2.5 副本健康自检:3 张系统表

2.5.1 system.replicas — 看每个副本的「健康度仪表盘」

sql
SELECT
    database, table,
    is_leader,                  -- 当前副本是否是 leader(决定 Merge 调度)
    is_readonly,                -- 是否只读(ZK 不通时会进只读保护)
    can_become_leader,
    parts_to_check,             -- 待校验 Part 数(>0 说明同步异常)
    queue_size,                 -- 待执行任务数(>200 报警)
    inserts_in_queue,           -- 待应用的 INSERT 数
    merges_in_queue,            -- 待应用的 MERGE 数
    log_max_index,              -- ZK 中最新的 log 序号
    log_pointer,                -- 本副本已消费到的 log 序号
    log_max_index - log_pointer AS lag,  -- 滞后多少
    absolute_delay,             -- 距离最新写入落后的秒数
    total_replicas,
    active_replicas
FROM system.replicas
FORMAT Vertical;

怎么看

  • is_readonly = 1立刻去看 Keeper / ZK 是否健康
  • lag > 1000 → 副本同步追不上写入速度,可能是网络瓶颈或对端机器负载高。
  • absolute_delay > 60s → 这个副本的数据已经「过时」了,读到它的查询会看到老数据。

2.5.2 system.replication_queue — 看「待执行队列」

sql
SELECT
    database, table,
    type,                       -- GET_PART / MERGE_PARTS / MUTATE_PART
    new_part_name,
    create_time,
    num_tries,                  -- 重试次数
    last_exception,             -- 失败原因
    postpone_reason             -- 推迟原因
FROM system.replication_queue
WHERE num_tries > 3 OR last_exception != ''
FORMAT Vertical;

如果这表里堆积大量任务,常见原因:

  • 磁盘满No space left on device —— 加盘或清老分区。
  • 对端 Part 不存在:源副本已经把那个 Part Merge 掉了,本副本拉不到 → 一般会自动恢复(从更新的 Part 拉)。
  • Schema 不一致:本节点的表结构和源副本不一样(手工改过 DDL)。

2.5.3 system.zookeeper — 直接看 ZK 节点

sql
SELECT name, value, ctime, mtime
FROM system.zookeeper
WHERE path = '/clickhouse/tables/01/events_local'
FORMAT Vertical;

可以列出该表在 ZK 里的所有子节点(log/replicas/mutations/block_numbers/ 等),是排查协调层问题的「终极武器」。

2.6 ⚠ 副本不是分片!再强调一次

写到这里如果你脑子里还在打架,做一个简单测试:

3 副本集群,3 台机器各 1TB 磁盘。能装多少业务数据?

答:1 TB。因为 3 副本是同一份数据印 3 份。

3 分片集群,3 台机器各 1TB 磁盘。能装多少业务数据?

答:3 TB。分片是切成 3 份。

3 分片 × 2 副本(共 6 台机器,每台 1TB)。能装多少业务数据?

答:3 TB。分片决定容量,副本决定可用性,副本不增容量。


3. Distributed 引擎:分布式查询入口

3.1 一句话定义

Distributed 引擎是「只负责路由、不存数据的虚拟表」。它知道集群里有几个分片、每个分片在哪台机器、查询时把 SQL 拆成 N 份发给各分片、写入时按分片键决定数据写到哪台机器。

3.2 「公司前台」类比

你(用户)                                    各部门(分片)
   │                                            ┌────┐ ┌────┐ ┌────┐
   │  「我要看 4 月份每个部门的销售额」           │ 销售│ │市场│ │财务│
   │       ↓                                    └────┘ └────┘ └────┘
┌──────────┐                                       ▲    ▲    ▲
│ 公司前台 │ ─── 拆成 3 份请求 ──────────────────►│    │    │
│Distributed│                                       │    │    │
│  路由表   │ ◄─── 收到 3 份小报表 ────────────────┘    │    │
└──────────┘ ◄─────────────────────────────────────────┘    │
   │  在前台合并成大报表                                       │
   │  (sum / 排序 / TopN)                          ◄──────────┘

返给你

关键洞察:

  • 前台自己什么数据都没有(Distributed 表本身不存数据);
  • 前台只知道部门通讯录remote_servers 配置);
  • 前台知道按什么把请求拆给谁sharding_key);
  • 前台负责「分发请求 + 合并结果」(fan-out / fan-in)。

3.3 本地表 + 分布式表的「黄金搭档」模式

生产里几乎所有 CK 集群都长这样:

每个节点上都有一个「本地表」 events_local(ReplicatedMergeTree,存数据)
每个节点上再建一个「分布式表」 events_all(Distributed,路由用)


          应用读 / 写都走 events_all,看不见 events_local
sql
-- 第 2 章的本地表(前面已经建过)
-- ENGINE = ReplicatedMergeTree(...)

-- 在每个节点上建 Distributed 视图,所有节点共用同一个名字
CREATE TABLE learn_ck.events_all ON CLUSTER ck_cluster
AS learn_ck.events_local                            -- 复用本地表 schema
ENGINE = Distributed(
    'ck_cluster',                                   -- 集群名(remote_servers 中的)
    'learn_ck',                                     -- 数据库
    'events_local',                                 -- 真实存数据的本地表
    cityHash64(user_id)                             -- 分片键:保证同一用户落同一分片
);

四个参数依次是:集群名、库、表、分片键。其中分片键决定写入时如何挑分片:

写入路径:
  INSERT INTO events_all (user_id=12345, ...) →
    Distributed 计算 cityHash64(12345) % 分片数 →
      落到 shard 2 →
        进 shard 2 中某个 replica 的 events_local →
          ReplicatedMergeTree 自动同步给同 shard 的其他 replica

3.4 分片键怎么选?数据倾斜怎么避免?

候选分片键优点缺点 / 风险
rand()数据绝对均匀同一用户的数据被打散到所有分片 → JOIN / 聚合时跨分片传输大
cityHash64(user_id)同用户落同分片 → 单用户聚合查询快有少量大用户时会数据倾斜
intDiv(toUnixTimestamp(ts), 86400 * 30)按月分片,方便冷热分层当月分片是「热分片」,写入压力集中
xxHash64(tenant_id)SaaS 多租户场景下保证租户隔离大租户也会倾斜

📌 生活类比:分片键就像快递分拣中心的「按邮编分拣」。如果按收件人姓名第一个字分,姓「张」的全堆一个柜子(倾斜);按邮编分则更均匀。

生产经验

  1. 高基数(user_id / device_id / order_id)+ Hash 函数是最常见的选择。
  2. 避免用低基数列(性别、城市、状态)当分片键 —— 几个值就把数据全压到几个分片上。
  3. 避免用时间字段当分片键 —— 当天写入只压在一个分片。
  4. 如果业务确实需要按时间分片,请改用「分区 + 多分片各持有所有时间」的组合:分片用 user_id Hash,分区用 toYYYYMM(ts),两者互补。

3.5 internal_replication:写入复制路径的关键开关

这是 ClickHouse 集群配置里最容易踩坑的参数。

xml
<remote_servers>
    <ck_cluster>
        <shard>
            <internal_replication>true</internal_replication>   <!-- ⭐ -->
            <replica><host>n1</host><port>9000</port></replica>
            <replica><host>n2</host><port>9000</port></replica>
        </shard>
        ...
    </ck_cluster>
</remote_servers>
取值行为推荐场景
true强烈推荐Distributed 写入时只写到分片中的一个副本,剩下的副本由 ReplicatedMergeTree 自己同步本地表用 ReplicatedMergeTree(生产标配)
falseDistributed 写入时自己向每个副本各写一份本地表用普通 MergeTree(无副本同步能力,不推荐)

为什么推荐 true

internal_replication = true:
  Distributed 写入 1 次 → ReplicatedMergeTree 通过 ZK 同步到其他副本
  优点:① 写入流量只 ×1;② 一致性由 ZK 保证;③ 容错好(一个副本写失败可重试)

internal_replication = false:
  Distributed 写入 N 次(N = 副本数)
  优点:不依赖 ZK
  缺点:① 写入流量 ×N;② 副本数据可能不一致(其中一次写失败就分裂);
        ③ 不能与 ReplicatedMergeTree 一起用(会双重写入)

记忆口诀:用了 Replicated*MergeTree,必开 internal_replication=true

3.6 分布式查询的 fan-out / fan-in

┌──────────────────────────────────────────────────────────────┐
│  SELECT event_type, count() FROM events_all                  │
│  WHERE event_time >= '2026-04-01' GROUP BY event_type;       │
└──────────────────────────────────────────────────────────────┘

                       Initiator 节点(接收 SQL 的节点)

            fan-out:把 SQL 改写后发给每个分片的某一个副本
                  │              │              │
                  ▼              ▼              ▼
           ┌─────────────┐┌─────────────┐┌─────────────┐
           │ shard1 副本 ││ shard2 副本 ││ shard3 副本 │
           │ 在本地 events_local 上跑              │
           │ SELECT event_type, count() ...       │
           │ GROUP BY event_type                  │
           │ → 返回中间聚合结果                     │
           └─────────────┘└─────────────┘└─────────────┘
                  │              │              │
                  └──────┬───────┴──────┬───────┘
                         ▼              ▼
                  fan-in:Initiator 收到 3 份小聚合

                  在本地做最终 merge:
                  countMerge() / sumMerge() / quantilesMerge() 等

                    最终结果返给客户端

为什么这个流程能撑「秒级聚合亿级数据」?

  1. 每个分片只扫自己那 1/N 数据 → IO 并行 N 倍。
  2. 聚合是「可结合」的:count、sum、min、max、distinct(用 HyperLogLog 草图)都能在分片本地预聚合,然后在 Initiator 上合并。这种「state → merge」模式是 ClickHouse 聚合函数的核心设计。
  3. 网络传输的是聚合状态而非原始行:一次 quantile() 在分片端只回传几 KB 的草图,而不是几百万行原始数据。

3.7 GLOBAL IN / GLOBAL JOIN:解决 N×N 问题

普通 IN 子查询在分布式下会有性能陷阱:

sql
-- ⚠ 朴素写法:会变成 N 次嵌套查询(N 个分片,每个分片对自己的 users 查询)
SELECT count() FROM events_all
WHERE user_id IN (SELECT user_id FROM users_all WHERE level = 'vip');

实际执行:每个分片在执行外层 events_local 的查询时,都会再次发起一遍内层 users_all 的分布式查询。N 个分片 × N 次内层查询 = N² 次。集群越大越慢。

解决方案:GLOBAL IN / GLOBAL JOIN

sql
SELECT count() FROM events_all
WHERE user_id GLOBAL IN (
    SELECT user_id FROM users_all WHERE level = 'vip'
);

GLOBAL 的语义:内层子查询在 Initiator 上只执行一次,结果作为临时表广播给所有分片,外层 SQL 只在分片内查。复杂度从 N² 降到 N。

📌 代价:广播的临时表必须能装进每个分片的内存。所以 GLOBAL IN 适合「右表小、左表大」的场景。如果右表也很大,就要换成 Distributed JOIN 或预先把数据分布到一致的分片键上(co-located join)。


4. 集群配置详解:remote_servers

集群拓扑的「源头真理」就是 config.xml 里这一段。完整结构:

xml
<remote_servers>
    <ck_cluster>                                  <!-- 集群名,自己取 -->
        <secret>shared-secret-for-cross-node</secret>  <!-- 节点间认证密钥(可选) -->

        <shard>                                   <!-- 第 1 个分片 -->
            <weight>1</weight>                    <!-- 分片权重,影响 rand() 分布 -->
            <internal_replication>true</internal_replication>
            <replica>
                <host>ck-n1.internal</host>
                <port>9000</port>
                <user>default</user>
                <password></password>
            </replica>
            <replica>
                <host>ck-n2.internal</host>
                <port>9000</port>
                <user>default</user>
                <password></password>
            </replica>
        </shard>

        <shard>                                   <!-- 第 2 个分片 -->
            <weight>1</weight>
            <internal_replication>true</internal_replication>
            <replica><host>ck-n3.internal</host><port>9000</port></replica>
            <replica><host>ck-n4.internal</host><port>9000</port></replica>
        </shard>

        <shard>                                   <!-- 第 3 个分片 -->
            <weight>1</weight>
            <internal_replication>true</internal_replication>
            <replica><host>ck-n5.internal</host><port>9000</port></replica>
            <replica><host>ck-n6.internal</host><port>9000</port></replica>
        </shard>
    </ck_cluster>
</remote_servers>

每个字段的含义:

字段作用
<shard>一个分片。所有 <replica> 在同一个分片下表示彼此互为镜像
<weight>该分片的写入权重。weight=2 的分片会拿到 2 倍数据(仅在分片键为 rand() 时生效)
<internal_replication>见 3.5 节
<replica>一个具体节点

📌 运行时验证集群拓扑

sql
SELECT cluster, shard_num, replica_num, host_name, host_address, port
FROM system.clusters
WHERE cluster = 'ck_cluster'
ORDER BY shard_num, replica_num;

5. 数据再平衡(Resharding):CK 没自动 rebalance

5.1 痛点

ClickHouse 不会在「加分片」「减分片」时自动重新分布数据。

加了第 4 个分片之后,老数据仍然只在原 3 个分片上,新写入才会按 4 分片 Hash。这就有两个问题:

  1. 老数据查询不均衡:老数据集中在前 3 个分片,第 4 个分片查它时是空。
  2. 历史数据再分片得自己干

5.2 几种解法

方案 1:INSERT ... SELECT 大法(最常用)

sql
-- 在新集群上建一张同 schema 但分片键改了的 events_all
CREATE TABLE learn_ck.events_all_v2 ON CLUSTER ck_cluster_v2
AS learn_ck.events_local_v2
ENGINE = Distributed('ck_cluster_v2', 'learn_ck', 'events_local_v2', cityHash64(user_id));

-- 灌入老数据
INSERT INTO learn_ck.events_all_v2
SELECT * FROM remote('ck_cluster', learn_ck.events_all);

通过 remote() 表函数从老集群读,写到新集群的 Distributed 表,新集群的 Distributed 会按新分片数 Hash。

优点:纯 SQL,可控。
缺点:得双倍存储、需要停业务窗口。

方案 2:clickhouse-copier

CK 官方提供的批量复制工具。基于 ZK 协调任务、支持断点续传、按分区粒度复制。适合 TB 级以上的数据搬家。

bash
clickhouse-copier --task-path=/clickhouse/copier/task1 \
                  --config=/etc/clickhouse-copier/config.xml \
                  --base-dir=/var/lib/clickhouse-copier

clickhouse-copier 在新版本中已被标注 "deprecated",但生产里仍广泛使用。新选项是直接用 INSERT ... SELECT + remote() 或者第 17 章会讲的 BACKUP / RESTORE

方案 3:物化视图 + 双写过渡

让旧表和新表同时被写入(业务层双写或用 MV 触发),切换流量后下线旧表。零停机但工程复杂度高。

5.3 给「未来不踩坑」的设计建议

  1. 分片数预留余量:上线第一天就按未来 3-5 年容量规划分片数。少不要紧(可以预留空分片),多了改起来非常痛。
  2. 分片键尽量稳定:选业务上不会变的字段(user_id、device_id),不要选可变字段(如 current_status)。
  3. 业务读全部走 Distributed 表:不要让业务直连 events_local,将来扩缩容才有灵活性。

6. ClickHouse Keeper:替代 ZooKeeper

6.1 为什么要替代 ZooKeeper

ZooKeeper 是 ClickHouse 早期的副本协调依赖,但它有几个先天痛点:

痛点ZooKeeper 的问题
语言Java 实现,吃内存(多副本压力下经常 GC 抖动)
依赖需要单独部署 JVM 集群,运维多一套
性能高频小事务下吞吐受限,CK 大集群常成为瓶颈
配置复杂zoo.cfgmyidlog4j.properties 一堆
网络分区znode watch 数量大时容易丢通知

6.2 ClickHouse Keeper 怎么解

ClickHouse Keeper 是用 C++ 重写的、协议兼容 ZooKeeper 的协调服务。核心特点:

  1. 协议兼容:客户端代码不用改,CK 服务端把 ZK 的 IP 换成 Keeper 的 IP 即可。
  2. C++ 实现:内存占用降到 ZK 的 1/3 ~ 1/4,没有 GC 抖动。
  3. 可以内嵌:可以作为独立进程部署(clickhouse-keeper),也可以直接嵌进 clickhouse-server 进程(小集群省机器)。
  4. 基于 Raft 共识算法:相比 ZAB(ZooKeeper Atomic Broadcast)在写入吞吐和延迟上都有改善。
  5. 快照与日志压缩更激进:恢复速度快很多。

6.3 部署形态对比

┌────────────────────────────────────────────────────────────────┐
│  方案 A:独立 ZooKeeper 集群(老方案)                            │
│                                                                 │
│   3 台 ZK   ←   CK 集群(6 节点)                                │
│   3 台 JVM       依赖网络打通                                    │
│                                                                 │
│  缺点:多一层运维 + Java GC 噪音                                  │
├────────────────────────────────────────────────────────────────┤
│  方案 B:独立 ClickHouse Keeper 集群                             │
│                                                                 │
│   3 台 Keeper(C++) ← CK 集群(6 节点)                         │
│                                                                 │
│  优点:协议兼容,纯 C++,内存友好                                 │
├────────────────────────────────────────────────────────────────┤
│  方案 C:Keeper 内嵌进 ClickHouse Server(小集群)                │
│                                                                 │
│   CK1 + Keeper(同进程)                                         │
│   CK2 + Keeper                                                  │
│   CK3 + Keeper                                                  │
│                                                                 │
│  优点:3 台机器搞定一切                                           │
│  注意:Keeper 节点数应为 3、5、7 奇数(Raft 要求)                │
└────────────────────────────────────────────────────────────────┘

6.4 配置示例(独立 Keeper)

xml
<!-- keeper_config.xml -->
<clickhouse>
    <keeper_server>
        <tcp_port>9181</tcp_port>
        <server_id>1</server_id>
        <log_storage_path>/var/lib/clickhouse/coordination/log</log_storage_path>
        <snapshot_storage_path>/var/lib/clickhouse/coordination/snapshots</snapshot_storage_path>

        <coordination_settings>
            <operation_timeout_ms>10000</operation_timeout_ms>
            <session_timeout_ms>30000</session_timeout_ms>
            <raft_logs_level>information</raft_logs_level>
        </coordination_settings>

        <raft_configuration>
            <server><id>1</id><hostname>kp-1</hostname><port>9234</port></server>
            <server><id>2</id><hostname>kp-2</hostname><port>9234</port></server>
            <server><id>3</id><hostname>kp-3</hostname><port>9234</port></server>
        </raft_configuration>
    </keeper_server>
</clickhouse>

CK Server 端引用 Keeper:

xml
<zookeeper>
    <node><host>kp-1</host><port>9181</port></node>
    <node><host>kp-2</host><port>9181</port></node>
    <node><host>kp-3</host><port>9181</port></node>
</zookeeper>

注意标签名仍然叫 <zookeeper> —— 这是为了客户端兼容,Keeper 完全可以平替。

📌 新项目首选 Keeper:截至 2024+,Altinity / 字节 / 阿里云 PolarDB-X for ClickHouse 都已经默认推 Keeper。新部署直接用 Keeper,不用考虑 ZK。


7. 真实案例:3 分片 × 2 副本读写完整流程

我们以 learn_ck.events(3 分片 × 2 副本,分片键 cityHash64(user_id))为例,完整跑一遍写读流程。

7.1 写入完整链路

sql
-- 应用执行:
INSERT INTO learn_ck.events_all
VALUES
    ('2026-04-17 10:00:00', 12345, 'click',     '{"page":"home"}'),
    ('2026-04-17 10:00:01', 67890, 'view',      '{"page":"detail"}'),
    ('2026-04-17 10:00:02', 12345, 'add_cart',  '{"sku":"A001"}');

第一步:客户端连到任一节点(假设 N1),把 SQL 交给 N1 上的 Distributed 引擎。

第二步:Distributed 计算每行的分片:

cityHash64(12345) % 3 = 1  → shard 2(行 1、行 3)
cityHash64(67890) % 3 = 0  → shard 1(行 2)

第三步:Distributed 把行按分片分组,只发给每个分片的一个副本(因为 internal_replication=true):

→ shard 1 的某个 replica(假设 N1):写入 row 2
→ shard 2 的某个 replica(假设 N4):写入 row 1, row 3

第四步:每个目标副本本地把数据走 ReplicatedMergeTree 写入 events_local,落 Part:

N1: events_local 落 Part `20260417_15_15_0` (1 行)
N4: events_local 落 Part `20260417_22_22_0` (2 行)

第五步:N1 / N4 各自在 ZK 登记 log,对应分片的另一副本(N2 / N3)从 ZK 拉取并 fetch Part:

N2 ← N1 拉取 `20260417_15_15_0`
N3 ← N4 拉取 `20260417_22_22_0`

第六步:客户端收到 OK,整个 INSERT 完成。注意:默认情况下,CK 不会等所有副本都同步完成才返回 —— 想要强一致写入用 insert_quorum 设置:

sql
SET insert_quorum = 2;       -- 至少 2 个副本确认才返回
SET insert_quorum_timeout = 10000;

7.2 查询完整链路

sql
SELECT event_type, uniqExact(user_id) AS uv
FROM learn_ck.events_all
WHERE event_time >= '2026-04-17'
GROUP BY event_type;

第一步:客户端连到任一节点(假设 N5),N5 成为 Initiator。

第二步:N5 上的 Distributed 引擎把 SQL 改写:

sql
-- 改写后给每个分片发的子查询
SELECT event_type, uniqExactState(user_id) AS uv_state
FROM learn_ck.events_local
WHERE event_time >= '2026-04-17'
GROUP BY event_type;

注意聚合函数从 uniqExact 变成了 uniqExactState,返回的是「聚合状态」而非最终值。

第三步:N5 在每个分片中挑一个健康副本(默认按 load_balancing 设置:random / nearest_hostname / in_order / first_or_random),并发发出去:

N5 → shard 1 选 N1 副本
N5 → shard 2 选 N3 副本
N5 → shard 3 选 N5 自己(本节点)

第四步:3 个分片各自跑本地 SQL,回传中间状态:

shard 1 返回:[('click', state_a), ('view', state_b), ...]
shard 2 返回:[('click', state_c), ('view', state_d), ...]
shard 3 返回:[('click', state_e), ('view', state_f), ...]

第五步:N5 在本地把状态合并:

sql
-- 对应的 final 步骤
SELECT event_type, uniqExactMerge(uv_state) AS uv
FROM (UNION ALL of the 3 results)
GROUP BY event_type;

第六步:把最终结果返给客户端。

7.3 失败场景演示

场景行为
shard 2 的 N3 副本宕机Distributed 自动切到 N4 副本(fallback),查询不中断
shard 2 全挂(N3 + N4 同时挂)查询报错:There is no leader for table 或返回部分结果(取决于 skip_unavailable_shards 设置)
写入时 ZK 不可用副本进 readonly,写入直接失败 → 业务必须有重试
Initiator 节点 N5 宕机客户端重连到任一节点重发 SQL 即可(无主架构的好处)

8. 📌 与同类系统的对比

维度ClickHouseMySQL 主从PostgreSQL 流复制GreenplumApache Doris FE-BE
复制模型拉取式 + ZK 协调binlog 推送(异步 / 半同步)WAL 流式推送共享存储 + segmentBE 之间 raft + FE 元数据
协调依赖ZK / Keeper无(主从直连)无(主从直连)gpsync 工具FE 自带元数据 BDB-JE
副本数任意(推荐 2-3)通常 1-2通常 1-2(流复制 + cascade)默认 2默认 3
分片Distributed 表手动指定应用层 / Vitess / TiDBCitus 扩展内置(hash / range)内置 tablet
自动 rebalance❌ 没有❌ 没有❌ 没有✅ gpexpand✅ 内置 tablet 调度
强一致写insert_quorum 可选rpl_semi_sync_master 半同步synchronous_commit 同步共享存储天然一致quorum 写
DDL 分发ON CLUSTER + ZK主上执行,binlog 复制主上执行,DDL 也走 WAL内置元数据广播FE 元数据广播
故障切换客户端重试任一健康节点MHA / orchestrator 工具repmgr / Patroni自带切换FE master 选举
数据放置感知不知道(自己手动配)不知道不知道知道(segment ID)知道(BE tablet)

核心洞察

  • ClickHouse 的设计更偏向「无中心、最终一致」。这让它在大集群下扩展性极强,但代价是「自动化运维」要自己堆。
  • Doris / StarRocks 这一代国产 OLAP 把「自动 rebalance / 自动副本迁移」做成内置能力,以更高的元数据复杂度换运维便利。CK 的哲学是「我做好引擎,运维你来」。
  • MySQL / PG 的复制都是「主写从读、有主从概念」,CK 的 ReplicatedMergeTree 是对等的 multi-master:任一副本都能写、都能读,只要 ZK 在协调。

9. 本章小结

┌─────────────────────────────────────────────────────────────────┐
│                        本章核心要点                                │
├─────────────────────────────────────────────────────────────────┤
│                                                                  │
│  ① 副本 vs 分片                                                   │
│     • 副本(Replica)= 镜像(解决高可用)                          │
│     • 分片(Shard)   = 切割(解决扩容)                            │
│     • 生产标配:N 分片 × M 副本(M ≥ 2)                           │
│                                                                  │
│  ② ReplicatedMergeTree                                           │
│     • 拉取式同步:源副本登记 ZK → 其他副本拉 Part                   │
│     • {shard}/{replica} 宏 + ZK 路径决定「谁和谁是兄弟」            │
│     • system.replicas / replication_queue / zookeeper 三件套排查 │
│                                                                  │
│  ③ Distributed                                                   │
│     • 不存数据的「路由器」                                          │
│     • 分片键决定写入分布(cityHash64(uid) 是最常见的安全选择)       │
│     • internal_replication=true:与 Replicated 配合的黄金姿势      │
│     • fan-out → 各分片预聚合 → fan-in 合并状态                     │
│     • GLOBAL IN / JOIN 解决 N×N 嵌套                              │
│                                                                  │
│  ④ Keeper 替代 ZK                                                 │
│     • C++ 实现 / 协议兼容 / 可内嵌                                 │
│     • 新项目首选                                                   │
│                                                                  │
│  ⑤ 没有自动 rebalance                                             │
│     • 加 / 减分片要自己 INSERT...SELECT 或 clickhouse-copier      │
│     • 上线前预留好分片数                                            │
│                                                                  │
│  ⑥ 一句话送给面试官:                                              │
│     「ReplicatedMergeTree 解决数据冗余,Distributed 解决数据拆分;│
│       两者通过 ZK/Keeper 协调,靠分片键决定数据放置,               │
│       靠 ON CLUSTER 完成 DDL 分发,构成 ClickHouse 的分布式骨架。」│
│                                                                  │
└─────────────────────────────────────────────────────────────────┘

10. 面试高频题

Q1:ClickHouse 的副本和分片是什么关系?为什么生产里都要混用?

考察点:是否真的把两个概念分清,能不能讲清各自的目的。

标准答案

  1. 副本(Replica)= 数据镜像:同一份数据在多台机器上各存一份,互为备份。目的是高可用读扩展(读请求可以负载均衡到多副本)。
  2. 分片(Shard)= 数据切割:按某个分片键(如 cityHash64(user_id))把数据分成 N 份,分别放到 N 台机器。目的是扩容量写并行
  3. 生产标配是「N 分片 × M 副本」(N ≥ 2, M ≥ 2):既扩容量又高可用。例如 3×2 集群(6 台机器):能存 3 倍数据,丢任意一台机器仍能读写。
  4. 实现层面:副本由 ReplicatedMergeTree 引擎 + ZooKeeper/Keeper 协调;分片由 Distributed 引擎 + remote_servers 配置 + 分片键完成路由。

加分项

  • 能补一句「副本和分片是正交的两个维度」,可以独立调整。
  • 能给反例「3 副本不增容量、3 分片不高可用」,证明理解到位。
  • 能提到 insert_quorum 强一致写、load_balancing 副本选择策略等细节。

易错点

  • 别把「3 副本」说成「数据量 ×3」—— 副本是镜像,总数据量不变
  • 别说「分片是高可用方案」—— 分片是扩容量方案,单点故障会丢一片数据

Q2:ReplicatedMergeTree 副本同步流程是怎样的?为什么用 ZK?

考察点:理解「副本协调」的本质 —— 用一个第三方公告板来代替 N×N 的两两握手。

标准答案

  1. 每个 ReplicatedMergeTree 表在 ZooKeeper(或 Keeper)下注册一棵 znode 子树,路径形如 /clickhouse/tables/{shard}/<table_name>。子树里包含:
    • log/:写入 / Merge / Mutation 任务的全局有序日志(ZK sequential node)
    • replicas/<replica_name>/queue/:每个副本的待执行队列
    • block_numbers/:用于 INSERT 的去重 + 块号分配
  2. 写流程
    • 副本 A 收到 INSERT → 排序 + 压缩 + 落本地 Part → 在 log/ 创建 sequential node 登记「需要 GET_PART」
    • 副本 B 后台任务轮询 log/,把新条目复制到自己的 queue/,逐个执行:从 A 的 9009 端口走 HTTP 拉取 Part 文件夹 → 校验 checksum → 落本地
  3. Merge 流程:每个分片有 leader 副本(在 leader_election/ 下选举),leader 决定哪些 Part 合并,写 log/ 登记 MERGE_PARTS,其他副本同样从 queue 拉取并执行相同的 merge(每个副本本地真正做合并,而不是从 leader 拉合并后的 Part)。
  4. 为什么用 ZK / Keeper:副本之间不直接握手,所有协调都通过这块 znode 子树。优点是 N 个副本只需要 N 条 watch,不会形成 N² 复杂度;缺点是 ZK / Keeper 成为关键依赖,挂了所有 Replicated 表进只读。

加分项

  • 能解释 block_id 的作用:客户端重试时 INSERT 不会写重复数据。
  • 能提到 Keeper 替代 ZK 的优势(C++ / 内嵌 / 没 GC)。
  • 能提到 system.replicas.absolute_delay 是衡量副本滞后的关键指标。

易错点

  • 别说「副本之间互相 push」—— 是 拉取式
  • 别说「ZK 存数据」—— ZK 只存元信息(几百字节的 JSON),数据走 HTTP 走 9009。

Q3:Distributed 表是什么?它存数据吗?如何选分片键?

考察点:能不能讲清 Distributed 的「路由表」本质 + 分片键设计原则。

标准答案

  1. Distributed 表本身不存数据,它是一张「虚拟表 / 路由器」,对外提供分布式的写入与查询接口。建表语法:
    sql
    CREATE TABLE events_all AS events_local
    ENGINE = Distributed('cluster_name', 'db', 'events_local', sharding_key);
  2. 写入:按 sharding_key 计算分片号(hash(key) % shard_count),把数据路由到对应分片的某个副本(如果 internal_replication=true 只写一个副本,剩下的由 Replicated 同步)。
  3. 查询:fan-out 把改写后的子 SQL 发给每个分片的某个副本(按 load_balancing 选副本),fan-in 把各分片的中间聚合状态在 Initiator 上 merge 成最终结果。
  4. 分片键选择原则
    • 高基数(user_id / device_id)+ Hash 函数(cityHash64 / xxHash64)→ 数据均匀
    • 与查询模式对齐:如果查询常按 user_id 聚合,分片键也用 user_id Hash → 单 user 聚合可在分片本地完成,不需要 reshuffle
    • 避免低基数(性别 / 城市)→ 数据倾斜
    • 避免单调时间字段 → 当天写入压一个分片

加分项

  • 能提到 weight 权重(按比例分流)。
  • 能解释 co-located join:两张表用相同分片键 → JOIN 可在分片本地完成,不需要广播。
  • 能说出读写都走 Distributed 的好处 —— 业务层无感知扩缩容。

易错点

  • 别说「Distributed 表存数据」—— 它是路由器。
  • 别用 rand() 做分片键 —— 数据分布是均匀了,但任何按用户聚合的查询都得跨分片 reshuffle。

Q4:internal_replication 这个参数是干嘛的?为什么推荐 true?

考察点:能不能避开生产里最经典的副本翻倍写入坑。

标准答案

  1. internal_replication<remote_servers> 配置中每个 <shard> 下的开关,决定 Distributed 写入时如何处理同分片下的多副本。
  2. true(推荐):Distributed 只写到分片中的一个副本,剩余副本由 ReplicatedMergeTree 通过 ZK 自己同步。优点:① 写入流量不翻倍;② 一致性由 ZK 保证;③ 容错好。
  3. false:Distributed 自己向每个副本各写一份,等于数据被写了 N 次。这种模式只适用于本地表是普通 MergeTree(无副本同步)的场景。
  4. 黄金搭档ReplicatedMergeTree + internal_replication=true
  5. 错误组合ReplicatedMergeTree + internal_replication=false → 数据被写了双倍(Distributed 写一次 + Replicated 同步一次),还可能因为副本各自再去重而引发分裂。

加分项

  • 能补充:insert_distributed_sync 设置控制 Distributed 是否同步等待远端确认。
  • 能提到对应 system.distribution_queue 视图查看 Distributed 的远端转发队列。

易错点

  • 别说「false 更安全(写多份)」—— 反而会写重复

Q5:分布式查询的 fan-out / fan-in 流程是什么?为什么聚合能扩展到亿级?

考察点:理解 ClickHouse 「中间状态」聚合模型。

标准答案

  1. fan-out:Initiator 节点(接收 SQL 的节点)把 SQL 改写后并发发给每个分片的某个副本。改写时把 Distributed 表名换成 local 表名,并把聚合函数改成对应的 *State 形式(countcountStateuniqExactuniqExactStatequantilequantileState)。
  2. 分片本地执行:每个分片在本地 MergeTree 上执行子 SQL,返回中间聚合状态(一段二进制 blob)而非最终结果。例如 uniqExact 的状态是 HashSet 序列化,quantile 的状态是采样草图(reservoir sampling)。
  3. fan-in:Initiator 把所有分片的状态收集,调用 *Merge 把状态合并成最终值。
  4. 为什么能扩展到亿级
    • 每个分片只扫自己的 1/N 数据 → IO 并行 N 倍
    • 网络只传聚合状态(KB 级),不传原始行(GB 级)
    • 聚合函数都设计成「可结合」的,符合 MapReduce 原理

加分项

  • 能解释 AggregateFunction / *State / *Merge 是如何对应到第 6 章的 AggregatingMergeTree 的(同一套机制)。
  • 能提到 prefer_localhost_replica 设置:本地有副本时优先打本地,省一次 RPC。

易错点

  • 别说「分片各自计算最终结果再相加」—— 对 uniqExactquantile 这种聚合,简单相加是错的,必须传中间状态。

Q6:什么是 GLOBAL IN / GLOBAL JOIN?什么时候用?

考察点:理解分布式查询里子查询的「N²」陷阱。

标准答案

  1. 问题:在 Distributed 表上做 ... WHERE x IN (SELECT ... FROM another_distributed),朴素实现下每个分片在外层执行时都会再次发起一遍内层分布式查询,复杂度变成 N²。
  2. GLOBAL IN / GLOBAL JOIN 的语义:内层子查询只在 Initiator 上执行一次,结果作为临时表广播给所有分片,外层在分片本地执行。复杂度从 N² 降到 N。
  3. 代价:广播的临时表必须能装进每个分片的内存(默认上限 max_distributed_connectionsmax_bytes_in_distributed_join)。
  4. 替代方案
    • Co-located JOIN:左右表用相同分片键 → JOIN 可以在分片本地完成,不需要广播
    • Dictionary:把右表灌进字典 → 字典在每个节点内存里有副本,不需要广播
    • 预聚合 / 物化视图:把 JOIN 结果预先 ETL 到宽表

加分项

  • 能给具体场景:vip 用户行为统计 → 右表是 vip 用户列表(小),左表是 events(大)→ 用 GLOBAL IN。
  • 能提到 enable_optimize_predicate_expression 等优化开关。

易错点

  • 别说「GLOBAL JOIN 比普通 JOIN 慢」—— 在分布式场景下,GLOBAL JOIN 是正确写法,普通 JOIN 才会因 N² 而慢。

Q7:ClickHouse Keeper 和 ZooKeeper 区别是什么?什么时候选 Keeper?

考察点:是否跟得上 ClickHouse 近几年的演进。

标准答案

  1. 协议兼容:ClickHouse Keeper 实现了 ZooKeeper 客户端协议,CK Server 可以无缝平替(只改 IP 端口)。
  2. 实现语言:ZK 是 Java(依赖 JVM、有 GC 抖动),Keeper 是 C++(无 GC、内存占用约为 ZK 的 1/3)。
  3. 一致性算法:ZK 用 ZAB;Keeper 用 Raft(实现自 NuRaft 库)。
  4. 部署方式:ZK 必须独立部署 JVM 集群;Keeper 可以独立部署,也可以内嵌进 clickhouse-server 进程,小集群可省掉 3 台机器。
  5. 何时选:新项目无脑选 Keeper;老项目 ZK 集群健康可以继续用,等机会再迁。

加分项

  • 能提到 Keeper 的 Raft 节点数必须是奇数(3/5/7),且过半数节点存活才能写。
  • 能提到 ZK → Keeper 的迁移工具 clickhouse-keeper-converter(把 ZK snapshot 转成 Keeper 格式)。

易错点

  • 别说「Keeper 不能用于其他系统」—— 协议兼容意味着任何 ZK 客户端都能连。
  • 别说「Keeper 取代了 etcd」—— 它替代的是 ZooKeeper,etcd 是另一回事。

Q8:ClickHouse 集群没有自动 rebalance,怎么扩缩容?

考察点:是否懂 CK 的运维痛点 + 实战处理思路。

标准答案

  1. 现状:CK 的 Distributed 引擎只在写入时按当前分片数路由,老数据不会自动迁移到新分片。也没有像 Doris / StarRocks 那样的 tablet 自动调度器。
  2. 扩容时
    • 加机器到 remote_servers,新数据按新分片数 Hash → 新分片开始接收数据
    • 老数据仍在原分片上,不影响读(Distributed 仍能从所有分片读到全量),但长期会形成不均衡
  3. 实操方案
    • 方案 A:INSERT ... SELECT + remote():建一个新集群,从老集群灌一遍数据,业务切换。最常用。
    • 方案 B:clickhouse-copier:批量复制工具,支持断点续传,但已 deprecated。
    • 方案 C:BACKUP / RESTORE:备份恢复,新版本支持。
    • 方案 D:物化视图 + 双写过渡:复杂度高但零停机。
  4. 预防:上线时一次性规划好 3-5 年容量的分片数;按业务量预估好分片键;避免后期改动。

加分项

  • 能提到 OPTIMIZE TABLE ... ON CLUSTER 在迁移后强制 merge。
  • 能提到「虚拟分片」的设计模式:用一致性 Hash 或权重,让单个物理节点持有多个逻辑分片,扩容时只迁移部分逻辑分片到新节点。

易错点

  • 别说「CK 自带 rebalance」—— 没有。
  • 别说「直接 ALTER TABLE ... MOVE PARTITION 就行」—— MOVE PARTITION 是同节点的盘间迁移,不是跨节点。

📌 下一章预告:第 14 章我们走出 ClickHouse 的「自留地」,看它如何与 Kafka、MySQL、PostgreSQL、S3、HDFS 这些上下游系统打通,做实时摄入和数据联邦查询。

🎬 可视化演示

演示加载缓慢或样式异常?点此在新标签页打开 ↗

💻 示例代码

python
"""
第 13 章 · 副本与分布式 · cluster_play.py

演示内容:
  1. 通过 Distributed 表写入 / 查询
  2. 打印 system.clusters(看集群拓扑)
  3. 打印 system.replicas(看副本健康)
  4. 演示 GLOBAL IN / GLOBAL JOIN 的差异
  5. 演示 insert_quorum 强一致写

运行前提:
  - 默认连 127.0.0.1:8123,default 用户,无密码,库 learn_ck
  - 单机情况下 ck_cluster 可能不存在,脚本会先尝试,失败时自动降级到本地表演示

依赖:
  pip install clickhouse-connect
"""

from __future__ import annotations

import os
import sys
import time
from contextlib import contextmanager

try:
    import clickhouse_connect
except ImportError:
    print("请先 pip install clickhouse-connect")
    sys.exit(1)


CK_HOST = os.getenv("CK_HOST", "127.0.0.1")
CK_PORT = int(os.getenv("CK_PORT", "8123"))
CK_USER = os.getenv("CK_USER", "default")
CK_PWD = os.getenv("CK_PWD", "")
CK_DB = os.getenv("CK_DB", "learn_ck")
CK_CLUSTER = os.getenv("CK_CLUSTER", "ck_cluster")


def banner(title: str) -> None:
    print("\n" + "=" * 78)
    print(f"  {title}")
    print("=" * 78)


def kv(rows, headers=None):
    """打印一个简易表格"""
    if not rows:
        print("  (空结果)")
        return
    if headers:
        widths = [
            max(len(str(h)), max((len(str(r[i])) for r in rows), default=0))
            for i, h in enumerate(headers)
        ]
        print("  " + "  ".join(str(h).ljust(widths[i]) for i, h in enumerate(headers)))
        print("  " + "  ".join("-" * w for w in widths))
        for r in rows:
            print("  " + "  ".join(str(r[i]).ljust(widths[i]) for i in range(len(headers))))
    else:
        for r in rows:
            print("  " + "  ".join(str(c) for c in r))


@contextmanager
def safe_block(title: str):
    """保证示例失败时不中断后续"""
    try:
        yield
    except Exception as exc:  # noqa: BLE001
        print(f"  ⚠ [{title}] 跳过:{exc}")


def main() -> None:
    print(f"连接 {CK_HOST}:{CK_PORT}(user={CK_USER}, db={CK_DB})...")
    client = clickhouse_connect.get_client(
        host=CK_HOST, port=CK_PORT, username=CK_USER, password=CK_PWD,
        database="default",
    )
    version = client.query("SELECT version()").result_rows[0][0]
    print(f"OK · ClickHouse 版本 {version}")

    # 检测集群是否存在
    cluster_exists = False
    with safe_block("check cluster"):
        rows = client.query(
            f"SELECT count() FROM system.clusters WHERE cluster = '{CK_CLUSTER}'"
        ).result_rows
        cluster_exists = rows and rows[0][0] > 0

    on_cluster = f"ON CLUSTER {CK_CLUSTER}" if cluster_exists else ""
    if not cluster_exists:
        print(f"\n注意:集群 {CK_CLUSTER} 不存在,将以单机/伪分布式模式演示。")
        print("      演示中会自动改用普通 MergeTree 与单分片 Distributed。")

    # ===== 1. 准备 schema =====
    banner("1. 创建演示库与表")
    client.command(f"CREATE DATABASE IF NOT EXISTS {CK_DB} {on_cluster}")
    client.command(f"DROP TABLE IF EXISTS {CK_DB}.events_all   {on_cluster} SYNC")
    client.command(f"DROP TABLE IF EXISTS {CK_DB}.events_local {on_cluster} SYNC")

    if cluster_exists:
        local_engine = (
            "ENGINE = ReplicatedMergeTree("
            "'/clickhouse/tables/{shard}/events_local', '{replica}')"
        )
    else:
        local_engine = "ENGINE = MergeTree()"

    client.command(f"""
        CREATE TABLE {CK_DB}.events_local {on_cluster} (
            event_time DateTime,
            user_id    UInt64,
            event_type LowCardinality(String),
            page       LowCardinality(String),
            properties String,
            country    LowCardinality(String) DEFAULT 'CN'
        ) {local_engine}
        PARTITION BY toYYYYMM(event_time)
        ORDER BY (event_type, user_id, event_time)
    """)

    if cluster_exists:
        client.command(f"""
            CREATE TABLE {CK_DB}.events_all {on_cluster}
            AS {CK_DB}.events_local
            ENGINE = Distributed('{CK_CLUSTER}', '{CK_DB}', 'events_local',
                                 cityHash64(user_id))
        """)
        write_table = f"{CK_DB}.events_all"
    else:
        write_table = f"{CK_DB}.events_local"

    print(f"  写入将走表:{write_table}")

    # ===== 2. 写入演示数据 =====
    banner("2. 写入演示数据")
    rows = [
        ("2026-04-17 09:00:00", 1001, "view",     "home",     '{"src":"app"}',  "CN"),
        ("2026-04-17 09:00:01", 1002, "view",     "home",     '{"src":"web"}',  "US"),
        ("2026-04-17 09:00:02", 1001, "click",    "banner",   '{"id":"b01"}',   "CN"),
        ("2026-04-17 09:00:03", 1003, "view",     "detail",   '{"sku":"A001"}', "JP"),
        ("2026-04-17 09:00:04", 1002, "add_cart", "detail",   '{"sku":"A001"}', "US"),
        ("2026-04-17 09:00:05", 1004, "view",     "home",     '{"src":"app"}',  "CN"),
        ("2026-04-17 09:00:06", 1005, "view",     "home",     '{"src":"app"}',  "CN"),
        ("2026-04-17 09:00:07", 1003, "pay",      "checkout", '{"amount":99}',  "JP"),
        ("2026-04-17 09:00:08", 1006, "view",     "detail",   '{"sku":"B002"}', "DE"),
        ("2026-04-17 09:00:09", 1007, "click",    "banner",   '{"id":"b02"}',   "CN"),
    ]
    cols = ["event_time", "user_id", "event_type", "page", "properties", "country"]
    t0 = time.time()
    client.insert(write_table, rows, column_names=cols)
    print(f"  写入 {len(rows)} 行,耗时 {(time.time()-t0)*1000:.1f} ms")

    # ===== 3. 集群拓扑 =====
    banner("3. system.clusters 集群拓扑")
    res = client.query(f"""
        SELECT cluster, shard_num, replica_num, host_name, port, is_local
        FROM system.clusters
        WHERE cluster = '{CK_CLUSTER}'
        ORDER BY shard_num, replica_num
    """)
    if res.result_rows:
        kv(res.result_rows, headers=["cluster", "shard", "replica", "host", "port", "is_local"])
    else:
        print(f"  集群 {CK_CLUSTER} 不存在(单机模式正常现象)")
        all_clusters = client.query(
            "SELECT DISTINCT cluster FROM system.clusters ORDER BY cluster"
        ).result_rows
        print(f"  本机已知的集群名:{[r[0] for r in all_clusters]}")

    # ===== 4. 副本健康度 =====
    banner("4. system.replicas 副本健康度(仅 Replicated*MergeTree 才有)")
    res = client.query(f"""
        SELECT
            database, table, is_leader, is_readonly,
            queue_size, log_max_index, log_pointer,
            log_max_index - log_pointer AS lag,
            absolute_delay, total_replicas, active_replicas
        FROM system.replicas
        WHERE database = '{CK_DB}'
    """)
    if res.result_rows:
        kv(
            res.result_rows,
            headers=["db", "table", "leader", "readonly", "queue", "log_max",
                     "log_ptr", "lag", "abs_delay", "total", "active"],
        )
    else:
        print("  本表不是 Replicated*MergeTree 引擎(单机模式正常现象)")

    # ===== 5. 查询:fan-out / fan-in =====
    banner("5. 简单聚合 SELECT(在分布式表上跑会自动 fan-out / fan-in)")
    if cluster_exists:
        sql = f"""
            SELECT event_type, uniqExact(user_id) AS uv, count() AS pv
            FROM {CK_DB}.events_all
            GROUP BY event_type
            ORDER BY pv DESC
        """
    else:
        sql = f"""
            SELECT event_type, uniqExact(user_id) AS uv, count() AS pv
            FROM {CK_DB}.events_local
            GROUP BY event_type
            ORDER BY pv DESC
        """
    res = client.query(sql)
    kv(res.result_rows, headers=["event_type", "uv", "pv"])

    # ===== 6. GLOBAL IN 演示(必须有分布式表才有意义) =====
    banner("6. IN vs GLOBAL IN(分布式查询的 N² → N 优化)")
    if cluster_exists:
        # 普通 IN:每个分片会再次发起内层分布式查询
        t0 = time.time()
        r1 = client.query(f"""
            SELECT count() FROM {CK_DB}.events_all
            WHERE user_id IN (
                SELECT user_id FROM {CK_DB}.events_all WHERE event_type = 'pay'
            )
        """).result_rows[0][0]
        t1 = time.time()
        # GLOBAL IN:内层只跑一次,广播给各分片
        r2 = client.query(f"""
            SELECT count() FROM {CK_DB}.events_all
            WHERE user_id GLOBAL IN (
                SELECT user_id FROM {CK_DB}.events_all WHERE event_type = 'pay'
            )
        """).result_rows[0][0]
        t2 = time.time()
        print(f"  普通 IN     →  {r1} 行  耗时 {(t1-t0)*1000:.1f} ms")
        print(f"  GLOBAL IN  →  {r2} 行  耗时 {(t2-t1)*1000:.1f} ms")
        print("  集群越大,差异越明显(GLOBAL IN 复杂度从 N² 降到 N)")
    else:
        print("  单机模式下两者无差异,跳过")

    # ===== 7. insert_quorum 强一致写 =====
    banner("7. insert_quorum 强一致写演示")
    if cluster_exists:
        try:
            client.command(f"""
                INSERT INTO {CK_DB}.events_all
                SETTINGS insert_quorum = 2, insert_quorum_timeout = 10000
                VALUES ('2026-04-17 10:00:00', 9999, 'view', 'home', '{{}}', 'CN')
            """)
            print("  ✓ 强一致写成功(至少 2 副本 ack)")
        except Exception as exc:
            print(f"  ⚠ 强一致写失败:{exc}")
            print("    可能原因:副本数不足、Keeper 不通、超时太短")
    else:
        print("  单机模式下 insert_quorum 无意义,跳过")

    # ===== 8. 看一下数据落到哪些分片 =====
    banner("8. 看每个分片各装了多少行(观察分片键的均衡度)")
    if cluster_exists:
        try:
            res = client.query(f"""
                SELECT hostName() AS host, count() AS rows
                FROM clusterAllReplicas('{CK_CLUSTER}', {CK_DB}.events_local)
                GROUP BY host
                ORDER BY host
            """)
            kv(res.result_rows, headers=["host", "rows"])
        except Exception as exc:
            print(f"  ⚠ 跳过:{exc}")
    else:
        res = client.query(f"SELECT count() FROM {CK_DB}.events_local")
        print(f"  本机总行数:{res.result_rows[0][0]}")

    print("\n演示完成。")


if __name__ == "__main__":
    main()

cluster_play.py ↗