ClickHouse 分布式集群

ClickHouse 的分布式能力使其能够水平扩展到数十个节点处理 PB 级数据。本文详解分片策略、副本复制、ZooKeeper 协调、集群配置和数据分布机制。

1. 分布式 ClickHouse 的架构模型

ClickHouse 的分布式架构由三个核心概念构成:

  • 分片(Shard):数据水平切分,每个分片包含数据的一个子集。查询时,请求被分发到所有分片并行执行,结果合并后返回
  • 副本(Replica):每个分片可以有多个副本,提供高可用性和读负载均衡。副本之间通过 ReplicatedMergeTree 引擎保持数据同步
  • 分布式表(Distributed Table):逻辑表,不存储实际数据,作为代理将查询路由到各分片的本地表

集群的典型拓扑:

                    ┌─────────────┐
                    │ Distributed │
                    │   events    │
                    └──────┬──────┘
                           │
           ┌───────────────┴───────────────┐
           ▼                               ▼
    ┌─────────────┐                 ┌─────────────┐
    │  Shard 1    │                 │  Shard 2    │
    │             │                 │             │
    │ ┌─────────┐ │                 │ ┌─────────┐ │
    │ │Replica 1A│ │                 │ │Replica 2A│ │
    │ └─────────┘ │                 │ └─────────┘ │
    │ ┌─────────┐ │                 │ ┌─────────┐ │
    │ │Replica 1B│ │                 │ │Replica 2B│ │
    │ └─────────┘ │                 │ └─────────┘ │
    └─────────────┘                 └─────────────┘

2. 集群配置

2.1 config.xml 集群定义

<!-- /etc/clickhouse-server/config.d/cluster.xml -->
<clickhouse>
    <remote_servers>
        <my_cluster>
            <!-- 分片 1 -->
            <shard>
                <internal_replication>true</internal_replication>
                <replica>
                    <host>shard1-replica1</host>
                    <port>9000</port>
                    <user>default</user>
                    <password>password</password>
                </replica>
                <replica>
                    <host>shard1-replica2</host>
                    <port>9000</port>
                    <user>default</user>
                    <password>password</password>
                </replica>
            </shard>
            
            <!-- 分片 2 -->
            <shard>
                <internal_replication>true</internal_replication>
                <replica>
                    <host>shard2-replica1</host>
                    <port>9000</port>
                    <user>default</user>
                    <password>password</password>
                </replica>
                <replica>
                    <host>shard2-replica2</host>
                    <port>9000</port>
                    <user>default</user>
                    <password>password</password>
                </replica>
            </shard>
        </my_cluster>
    </remote_servers>
</clickhouse>

2.2 ZooKeeper / ClickHouse Keeper 配置

ReplicatedMergeTree 需要 ZooKeeper 或 ClickHouse Keeper 来协调副本之间的元数据:

<clickhouse>
    <zookeeper>
        <node>
            <host>zk1</host>
            <port>2181</port>
        </node>
        <node>
            <host>zk2</host>
            <port>2181</port>
        </node>
        <node>
            <host>zk3</host>
            <port>2181</port>
        </node>
    </zookeeper>
</clickhouse>

ClickHouse Keeper 是 ZooKeeper 的替代实现,与 ClickHouse 二进制集成,无需单独部署。

3. 数据分片策略

3.1 分片键选择

分片键决定数据写入哪个分片。好的分片键应该使数据均匀分布,同时让查询能定位到最少的分片。

-- 使用 rand() 均匀随机分片
CREATE TABLE events_distributed AS events_local
ENGINE = Distributed(my_cluster, default, events_local, rand());
-- 优点:分布绝对均匀
-- 缺点:范围查询需要扫描所有分片

-- 使用 user_id 的哈希(同一用户的数据在同一分片)
CREATE TABLE events_distributed AS events_local
ENGINE = Distributed(my_cluster, default, events_local, cityHash64(user_id));
-- 优点:按用户查询只需访问一个分片
-- 缺点:如果有热点用户,可能导致数据倾斜

-- 使用时间范围分片(不同时间段在不同分片)
CREATE TABLE events_distributed AS events_local
ENGINE = Distributed(my_cluster, default, events_local, toYYYYMMDD(event_time));
-- 优点:按时间范围查询只需访问部分分片
-- 缺点:新数据总是写入一个分片,写入不分散

3.2 分片策略对比

策略写入分布点查效率范围查询风险
rand()均匀需全部分片需全部分片无热点
hash(key)均匀单个分片需全部分片可能倾斜
time集中时间分区仅相关分片写入热点
geo地理分布地区查询快需全部分片地区大小不均

4. 分布式表的读写

4.1 写入分布式表

-- 写入分布式表,数据自动路由到正确的分片
INSERT INTO events_distributed (event_time, user_id, event_type)
VALUES (now(), 12345, 'click');

-- 但批量写入推荐直接写本地表(性能更好)
-- 每个应用实例写入特定分片
INSERT INTO events_local (event_time, user_id, event_type)
VALUES (now(), 12345, 'click');

4.2 查询分布式表

-- 查询分布式表(自动分发到所有分片,聚合结果)
SELECT 
    event_type,
    count() AS cnt,
    avg(value) AS avg_value
FROM events_distributed
WHERE event_time > today() - 7
GROUP BY event_type;

-- ClickHouse 自动:
-- 1. 将查询发送到所有分片的副本(负载均衡)
-- 2. 每个分片本地执行查询
-- 3. 收集所有分片结果
-- 4. 合并(执行最终的 GROUP BY、ORDER BY、LIMIT)
-- 5. 返回结果

-- 强制在所有分片上查询(即使可能有重复)
SELECT count() FROM events_distributed;

-- 查看分布式查询执行计划
EXPLAIN SELECT count() FROM events_distributed;

4.3 GLOBAL JOIN

在分布式查询中做 JOIN 需要特别注意:

-- ❌ 问题:每个分片只加载本分片的 orders,JOIN 结果不完整
SELECT e.*, o.amount
FROM events_distributed e
JOIN orders_distributed o ON e.order_id = o.order_id;

-- ✅ 方案 1:GLOBAL JOIN,复制右表到所有分片
SELECT e.*, o.amount
FROM events_distributed e
GLOBAL JOIN orders_distributed o ON e.order_id = o.order_id;

-- ✅ 方案 2:JOIN 引擎(小表)
-- 在每台服务器上创建包含完整数据的 JOIN 表
CREATE TABLE orders_all (...) ENGINE = Join(ANY, LEFT, order_id);

-- ✅ 方案 3:字典(小维表)
-- 使用字典加载小表到内存

5. 副本与高可用

5.1 ReplicatedMergeTree

-- 在每个分片的每个副本上创建本地表
CREATE TABLE events_local ON CLUSTER my_cluster (
    event_time DateTime,
    user_id UInt64,
    event_type String
) ENGINE = ReplicatedMergeTree(
    '/clickhouse/tables/{shard}/events',  -- ZooKeeper 路径
    '{replica}'                              -- 副本标识
)
PARTITION BY toYYYYMM(event_time)
ORDER BY (event_time, user_id);

宏变量(macros)在每个节点上不同:

<!-- shard1-replica1 的配置 -->
<macros>
    <shard>01</shard>
    <replica>shard1-replica1</replica>
</macros>

<!-- shard1-replica2 的配置 -->
<macros>
    <shard>01</shard>
    <replica>shard1-replica2</replica>
</macros>

<!-- shard2-replica1 的配置 -->
<macros>
    <shard>02</shard>
    <replica>shard2-replica1</replica>
</macros>

5.2 副本同步机制

写入流程:
1. 客户端写入副本 1A
2. 副本 1A 写入本地数据
3. 副本 1A 在 ZooKeeper 创建任务
4. 副本 1B 监听 ZK,获取复制任务
5. 副本 1B 从 1A 拉取新 part
6. 所有副本数据一致
-- 检查副本状态
SELECT
    database,
    table,
    is_leader,
    is_readonly,
    absolute_delay,
    queue_size
FROM system.replicas
WHERE table = 'events_local';

-- 查看副本延迟
SELECT
    database,
    table,
    replica_name,
    absolute_delay
FROM system.replicas
ORDER BY absolute_delay DESC;

6. ON CLUSTER 执行

-- 在所有分片/副本上同时执行 DDL
CREATE TABLE events_local ON CLUSTER my_cluster (
    event_time DateTime,
    user_id UInt64
) ENGINE = ReplicatedMergeTree(...)
ORDER BY event_time;

-- 删除集群上的表
DROP TABLE IF EXISTS events_local ON CLUSTER my_cluster;

-- 修改表结构(所有节点)
ALTER TABLE events_local ON CLUSTER my_cluster
ADD COLUMN event_subtype String AFTER event_type;

7. 数据迁移与再平衡

7.1 副本间数据同步

-- 强制同步所有副本
SYSTEM SYNC REPLICA events_local;

-- 停止副本接收
SYSTEM STOP REPLICATED SENDS;
SYSTEM START REPLICATED SENDS;

7.2 分片间数据迁移

-- 将某节点的数据复制到新分片
INSERT INTO events_local_new
SELECT * FROM remote('old-node', 'default', 'events_local');

8. 监控分布式集群

-- 查看集群节点状态
SELECT * FROM system.clusters WHERE cluster = 'my_cluster';

-- 分布式查询耗时分析
SELECT
    query,
    query_duration_ms,
    read_rows,
    read_bytes,
    result_rows,
    is_initial_query
FROM system.query_log
WHERE type = 'QueryFinish'
ORDER BY event_time DESC
LIMIT 10;

-- 监控分片间数据量差异
SELECT
    hostName() AS host,
    count() AS rows,
    formatReadableSize(sum(bytes)) AS size
FROM clusterAllReplicas('my_cluster', default.events_local)
GROUP BY host;

9. 总结

ClickHouse 分布式设计的核心要点:

层面关键决策建议
分片策略均匀性 vs 查询定位性根据查询模式选择
副本数可用性需求生产至少 2 副本
协调服务ZooKeeper vs KeeperKeeper 更轻量
写入方式分布式表 vs 本地表大批量写本地表
JOIN 策略GLOBAL JOIN vs 字典小表用字典/Join 引擎
数据倾斜监控各分片大小及时调整分片键

ClickHouse 的分布式架构使其能够线性扩展到数十个节点处理 PB 级数据,但成功的分布式部署需要仔细的分片设计和持续的集群监控。

继续阅读

探索更多技术文章

浏览归档,发现更多关于系统设计、工具链和工程实践的内容。

全部文章 返回首页

「数据库」更多文章

  1. ClickHouse 表引擎详解
  2. ClickHouse 监控与运维
  3. ClickHouse 生产案例与最佳实践