上一篇【第46篇】ClickHouse ReplicatedMergeTree原理深度解析
下一篇【第48篇】ClickHouse Distributed引擎原理_分布式读写核心流程


摘要

本文是《Clickhouse从入门到精通》系列博客的第47篇文章,深入讲解ClickHouse数据分片机制与集群配置的核心原理。数据分片是ClickHouse实现水平扩展的基石,通过将大规模数据集分散到多个节点上,实现存储和计算的横向扩展。本文将系统介绍集群配置语法(remote_servers)、internal_replication参数的核心作用、宏变量(macros)的高级用法、多种分片策略的对比分析、ON CLUSTER分布式DDL的执行机制与监控方法,以及分布式建表的最佳实践。文章包含大量实战SQL示例和完整的3分片2副本集群搭建案例,帮助读者深入理解ClickHouse分布式架构的设计哲学与工程实践。

关键词:ClickHouse、数据分片、集群配置、分布式DDL、宏变量


1. 引言

在大数据时代,单机数据库在面对PB级数据存储和高并发查询时往往力不从心。ClickHouse作为一款专为OLAP场景设计的列式数据库,其分布式架构设计是其核心竞争力之一。数据分片(Sharding)是ClickHouse实现水平扩展的核心机制,它允许将一张逻辑表的数据按照某种规则分散存储到多个节点上,每个节点只存储部分数据,从而实现存储和计算的线性扩展。

与许多分布式系统不同,ClickHouse的分片机制设计哲学是"简单而灵活"——它不强制使用特定的分片策略,而是将分片规则的制定权交给用户,同时通过Distributed引擎和ON CLUSTER DDL等机制,让分布式操作尽可能对用户透明。

理解ClickHouse的数据分片机制,需要掌握以下几个核心概念:

  1. 集群(Cluster):由多个分片(Shard)组成的逻辑集合,每个分片可以包含多个副本(Replica)。
  2. 分片(Shard):数据水平分割的基本单位,每个分片存储整个数据集的一个子集。
  3. 副本(Replica):分片的数据冗余备份,用于高可用和读负载均衡。
  4. 分布式表(Distributed Table):一种特殊的逻辑表引擎,本身不存储数据,而是作为数据分片的路由层。

本文将围绕这些核心概念,从集群配置入手,逐步深入分片策略、分布式DDL、建表最佳实践等关键主题。


2. 集群配置详解

2.1 remote_servers配置结构

ClickHouse的集群定义在配置文件 config.xml(或 config.d/ 目录下的独立配置文件)中,通过 <remote_servers> 标签进行配置。一个完整的集群定义包含集群名称、分片列表、每个分片的副本列表等信息。

基础配置语法
<clickhouse>
    <remote_servers>
        <my_cluster>  <!-- 集群名称,在ON CLUSTER子句和Distributed引擎中引用 -->
            <shard>   <!-- 分片1 -->
                <replica>  <!-- 分片1的副本1 -->
                    <host>node1.example.com</host>
                    <port>9000</port>
                    <user>default</user>
                    <password>secret</password>
                </replica>
                <replica>  <!-- 分片1的副本2 -->
                    <host>node2.example.com</host>
                    <port>9000</port>
                    <user>default</user>
                    <password>secret</password>
                </replica>
            </shard>
            <shard>   <!-- 分片2 -->
                <replica>
                    <host>node3.example.com</host>
                    <port>9000</port>
                </replica>
            </shard>
        </my_cluster>
    </remote_servers>
</clickhouse>
配置参数详解
标签/参数必填说明
<remote_servers>所有集群定义的容器
<cluster_name>集群名称,全局唯一
<shard>分片定义,可包含多个
<replica>副本定义,每个分片至少一个
<host>节点主机名或IP地址
<port>ClickHouse TCP端口(默认9000)
<user>连接用户名(默认default)
<password>连接密码
<secure>是否使用TLS加密连接(0/1)
<compression>是否启用压缩(0/1)
internal_replication参数深度解析

<internal_replication> 是集群配置中最重要的参数之一,它决定了数据写入时分片内多个副本之间的数据同步方式。

<shard>
    <internal_replication>true</internal_replication>
    <replica>...</replica>
    <replica>...</replica>
</shard>

internal_replication = true(推荐)

当设置为 true 时,ClickHouse只会将数据写入分片的一个副本,然后依靠该副本自身的复制机制(ReplicatedMergeTree引擎)将数据同步到其他副本。

工作流程

  1. 写入请求到达Distributed表
  2. Distributed表根据sharding_key计算目标分片
  3. 数据只写入该分片的第一个健康副本
  4. 该副本通过ZooKeeper协调,将数据异步同步到其他副本

优点

  • 避免重复写入,减少网络开销
  • 利用ReplicatedMergeTree原生的复制机制,保证数据一致性
  • 写入性能更好

缺点

  • 依赖ZooKeeper,ZK压力较大时需要关注

internal_replication = false

当设置为 false 时,ClickHouse会将数据同时写入分片内的所有副本。

工作流程

  1. 写入请求到达Distributed表
  2. 数据被同时发送到该分片的所有副本
  3. 各副本独立写入,不通过ZooKeeper同步

优点

  • 不依赖ZooKeeper,适合ZK不稳定的场景
  • 写入路径简单

缺点

  • 网络带宽是true模式的倍数(副本数个)
  • 可能出现数据不一致(部分副本写入失败)
  • 不推荐在生产环境使用
对比总结
特性internal_replication=trueinternal_replication=false
写入目标仅第一个健康副本所有副本
副本同步方式ReplicatedMergeTree自动同步客户端同时写入
网络开销低(只写一份)高(写N份)
数据一致性强一致(依赖ZK)最终一致(可能不一致)
ZooKeeper依赖
生产推荐度⭐⭐⭐⭐⭐⭐⭐

2.2 多集群配置示例

在实际生产环境中,通常需要根据业务特点配置多个集群。例如,可以配置一个用于在线服务的集群和一个用于离线分析的集群。

<remote_servers>
    <!-- 在线服务集群:3分片2副本 -->
    <online_cluster>
        <shard>
            <internal_replication>true</internal_replication>
            <replica><host>ch-online-01</host><port>9000</port></replica>
            <replica><host>ch-online-02</host><port>9000</port></replica>
        </shard>
        <shard>
            <internal_replication>true</internal_replication>
            <replica><host>ch-online-03</host><port>9000</port></replica>
            <replica><host>ch-online-04</host><port>9000</port></replica>
        </shard>
        <shard>
            <internal_replication>true</internal_replication>
            <replica><host>ch-online-05</host><port>9000</port></replica>
            <replica><host>ch-online-06</host><port>9000</port></replica>
        </shard>
    </online_cluster>

    <!-- 离线分析集群:2分片1副本(大存储节点) -->
    <offline_cluster>
        <shard>
            <internal_replication>false</internal_replication>
            <replica><host>ch-offline-01</host><port>9000</port></replica>
        </shard>
        <shard>
            <internal_replication>false</internal_replication>
            <replica><host>ch-offline-02</host><port>9000</port></replica>
        </shard>
    </offline_cluster>

    <!-- 跨机房集群:每个分片跨越两个机房 -->
    <cross_dc_cluster>
        <shard>
            <internal_replication>true</internal_replication>
            <!-- 机房A的副本 -->
            <replica><host>ch-dc1-01</host><port>9000</port></replica>
            <!-- 机房B的副本 -->
            <replica><host>ch-dc2-01</host><port>9000</port></replica>
        </shard>
    </cross_dc_cluster>
</remote_servers>

2.3 宏变量(macros)配置与使用

在分布式集群中,每个节点都需要知道自己的"身份"——它是哪个分片、哪个副本。ClickHouse通过**宏变量(macros)**机制来解决这个问题。宏变量在配置文件中定义,然后在CREATE TABLE等DDL语句中通过 {macro_name} 语法引用。

配置宏变量

在每个节点的 config.xmlconfig.d/macros.xml 中配置:

<!-- 节点 ch-online-01 的配置 -->
<clickhouse>
    <macros>
        <cluster>online_cluster</cluster>
        <shard>01</shard>
        <replica>ch-online-01</replica>
    </macros>
</clickhouse>
<!-- 节点 ch-online-02 的配置 -->
<clickhouse>
    <macros>
        <cluster>online_cluster</cluster>
        <shard>01</shard>
        <replica>ch-online-02</replica>
    </macros>
</clickhouse>
在DDL中使用宏变量

配置好宏变量后,可以使用统一的DDL语句在所有节点上创建本地表:

-- 在所有节点上执行相同的DDL(通过ON CLUSTER自动分发)
CREATE TABLE IF NOT EXISTS local_user_events ON CLUSTER '{cluster}'
(
    event_date Date,
    user_id UInt64,
    event_type String,
    event_value Float64
)
ENGINE = ReplicatedMergeTree('/clickhouse/tables/{shard}/user_events', '{replica}')
PARTITION BY toYYYYMM(event_date)
ORDER BY (user_id, event_date);

在这个例子中:

  • {cluster} 被替换为 online_cluster
  • {shard} 被替换为各节点自身的分片编号(如 010203
  • {replica} 被替换为各节点自身的副本名称(如 ch-online-01
常用宏变量命名规范
宏变量名含义示例值
{cluster}集群名称online_cluster
{shard}分片编号010203
{replica}副本标识ch-node-01replica2
{database}数据库名analytics
{table}表名user_events
宏变量的高级用法

宏变量不仅可以用于ReplicatedMergeTree的ZK路径,还可以用于:

  1. Kafka引擎表的消费者组名称
CREATE TABLE kafka_events ON CLUSTER '{cluster}'
ENGINE = Kafka()
SETTINGS
    kafka_broker_list = 'kafka-01:9092,kafka-02:9092',
    kafka_topic_list = 'events',
    kafka_group_name = 'ch_consumer_{shard}',  -- 每个分片独立的消费者组
    kafka_format = 'JSONEachRow';
  1. 字典文件的路径
<dictionaries>
    <dictionary>
        <path>/data/dictionaries/{cluster}/{shard}/geo_dict.txt</path>
    </dictionary>
</dictionaries>

3. 查看集群信息

3.1 system.clusters系统表

配置好集群后,可以通过 system.clusters 系统表查看集群的详细拓扑信息。

SELECT
    cluster,
    shard_num,
    shard_weight,
    replica_num,
    host_name,
    port,
    user,
    is_local
FROM system.clusters
ORDER BY cluster, shard_num, replica_num;

字段说明

字段类型说明
clusterString集群名称
shard_numUInt32分片编号(从1开始)
shard_weightUInt32分片权重(用于加权分片)
replica_numUInt32副本编号(从1开始)
host_nameString主机名
host_addressStringIP地址
portUInt16端口号
userString用户名
is_localUInt8是否为本地节点
errors_countUInt32连接错误次数
estimated_recovery_timeUInt32预计恢复时间(秒)

示例输出

cluster         | shard_num | shard_weight | replica_num | host_name      | port | user    | is_local
----------------+-----------+--------------+-------------+----------------+------+---------+----------
online_cluster  |         1 |            1 |           1 | ch-online-01   | 9000 | default |        1
online_cluster  |         1 |            1 |           2 | ch-online-02   | 9000 | default |        0
online_cluster  |         2 |            1 |           1 | ch-online-03   | 9000 | default |        0
...

3.2 cluster()表函数

cluster() 表函数允许直接查询远程集群上的表,无需创建Distributed表。

语法

cluster(cluster_name, database, table, sharding_key)

示例

-- 查询远程集群上各分片的数据量
SELECT
    shard_num,
    count() AS rows
FROM cluster('online_cluster', 'default', 'local_user_events', rand())
GROUP BY shard_num
ORDER BY shard_num;

注意cluster() 表函数每次调用都会建立新的连接,不适合高频使用。生产环境应使用Distributed表引擎。


4. 分片策略详解

分片策略决定了数据如何分布到各个分片。ClickHouse通过 sharding_key 参数来指定分片规则。合理的分片策略应该保证数据均匀分布,同时兼顾查询效率。

4.1 随机分片(rand())

随机分片使用 rand() 函数作为sharding_key,数据会被随机分配到各个分片。

CREATE TABLE distributed_user_events ON CLUSTER 'online_cluster'
ENGINE = Distributed('online_cluster', 'default', 'local_user_events', rand())

特点

  • 数据分布均匀性好(大数定律)
  • 实现简单
  • 无法保证相同业务键的数据在同一分片

适用场景

  • 数据没有自然的分片键
  • 查询主要是全表扫描或跨分片聚合
  • 不需要按照某个业务字段进行分片级剪枝

4.2 哈希分片(cityHash64())

哈希分片通过对业务键计算哈希值,然后将哈希值对分片数取模,决定数据所属的分片。

-- 按照user_id进行哈希分片
CREATE TABLE distributed_user_events ON CLUSTER 'online_cluster'
ENGINE = Distributed('online_cluster', 'default', 'local_user_events', cityHash64(user_id))

常用哈希函数对比

函数分布均匀性性能说明
cityHash64()⭐⭐⭐⭐⭐⭐⭐⭐⭐推荐,分布最均匀
sipHash64()⭐⭐⭐⭐⭐⭐⭐加密级哈希,更慢但更安全
halfMD5()⭐⭐⭐⭐⭐兼容性较好
murmurHash3_64()⭐⭐⭐⭐⭐⭐⭐⭐性能与均匀性俱佳

特点

  • 相同 user_id 的数据总是落在同一分片
  • 支持分片级查询剪枝(where条件包含分片键时)
  • 哈希冲突可能导致轻微的数据倾斜

适用场景

  • 需要按照某个业务维度(如user_id、order_id)进行数据局部性聚合
  • 查询条件中经常包含分片键

4.3 自定义分片键表达式

除了使用简单的哈希函数,还可以使用复杂的表达式作为sharding_key。

-- 按照日期+用户ID的复合分片键
ENGINE = Distributed(
    'online_cluster',
    'default',
    'local_user_events',
    cityHash64(toString(event_date) || '_' || toString(user_id))
)
-- 按照范围分片(假设分片1存1月数据,分片2存2月数据...)
ENGINE = Distributed(
    'online_cluster',
    'default',
    'local_user_events',
    toMonth(event_date) % 3 + 1  -- 3个分片
)

警告:自定义分片键表达式必须是确定性的(相同输入产生相同输出),否则可能导致数据不一致。

4.4 各分片策略对比

分片策略数据均匀性查询剪枝相同键同分片实现复杂度推荐指数
随机分片(rand())⭐⭐⭐⭐⭐⭐⭐
哈希分片(cityHash64)⭐⭐⭐⭐⭐⭐⭐⭐⭐⭐⭐⭐
取模分片(% shard_count)⭐⭐⭐⭐⭐⭐
范围分片⭐⭐✅✅取决于范围⭐⭐⭐⭐⭐
一致性哈希⭐⭐⭐⭐⭐⭐⭐⭐⭐⭐⭐⭐

数据倾斜风险分析

  1. 随机分片:理论上无数据倾斜,但实际中由于随机数分布特性,小数据量时可能有偏差。
  2. 哈希分片:如果分片键的基数低(如只有几个不同值),会导致严重的数据倾斜。
  3. 取模分片:当分片数变化时,大部分数据需要重新分片。
  4. 范围分片:数据分布最不均匀,容易出现热点分片。

5. ON CLUSTER分布式DDL

5.1 语法与使用

ON CLUSTER 子句允许在集群的所有节点上执行相同的DDL语句,极大简化了分布式环境下的运维操作。

支持的DDL类型

-- 1. 创建本地表(在每个节点上创建相同的本地表)
CREATE TABLE local_user_events ON CLUSTER 'online_cluster'
(
    event_date Date,
    user_id UInt64,
    event_type String,
    event_value Float64
)
ENGINE = ReplicatedMergeTree('/clickhouse/tables/{shard}/user_events', '{replica}')
PARTITION BY toYYYYMM(event_date)
ORDER BY (user_id, event_date);

-- 2. 创建分布式表(在每个节点上创建指向本地表的分布式视图)
CREATE TABLE distributed_user_events ON CLUSTER 'online_cluster'
ENGINE = Distributed('online_cluster', 'default', 'local_user_events', cityHash64(user_id));

-- 3. 添加列
ALTER TABLE local_user_events ON CLUSTER 'online_cluster'
ADD COLUMN IF NOT EXISTS event_source String DEFAULT 'unknown';

-- 4. 删除表
DROP TABLE IF EXISTS local_user_events ON CLUSTER 'online_cluster';

-- 5. 创建数据库
CREATE DATABASE IF NOT EXISTS analytics ON CLUSTER 'online_cluster';

-- 6. 创建物化视图
CREATE MATERIALIZED VIEW mv_user_daily ON CLUSTER 'online_cluster'
ENGINE = ReplicatedAggregatingMergeTree(...)
...

5.2 DDL执行机制

当执行一条 ON CLUSTER DDL时,ClickHouse内部的执行流程如下:

客户端
  │
  ▼
执行节点(接收DDL的节点)
  │
  │  1. 将DDL任务写入ZooKeeper路径:
  │     /clickhouse/task_queue/ddl/{cluster}/query-[uuid]
  │
  ▼
ZooKeeper(DDL任务队列)
  ▲
  │  2. 各节点监听ZK路径,发现新任务
  │
  │
各节点(分片的所有副本)
  │
  │  3. 拉取任务,在本节点执行DDL
  │
  │  4. 将执行结果写回ZK:
  │     /clickhouse/task_queue/ddl/{cluster}/query-[uuid]/finished/[node_id]
  │
  ▼
执行节点
  │
  │  5. 收集所有节点的执行结果
  │
  ▼
客户端(返回执行结果)

关键ZooKeeper路径

/clickhouse/task_queue/ddl/
└── online_cluster/
    ├── query-abc123def456     # DDL任务节点
    │   ├── query      # 存储DDL语句
    │   ├── hosts      # 目标主机列表
    │   ├── finished/  # 各节点执行结果
    │   │   ├── node1 -> "success"
    │   │   ├── node2 -> "success"
    │   │   └── node3 -> "error: ..."
    │   └── active    # 当前正在执行的节点
    └── ...

5.3 超时与重试配置

分布式DDL的默认超时时间是 180秒,可以通过以下方式调整:

-- 方式1:会话级别设置
SET distributed_ddl_task_timeout = 300;  -- 单位:秒

-- 方式2:语句级别设置
CREATE TABLE ... ON CLUSTER 'online_cluster'
SETTINGS distributed_ddl_task_timeout = 300;

相关配置参数

参数默认值说明
distributed_ddl_task_timeout180DDL任务等待所有节点完成的超时时间
distributed_ddl_task_max_tries3执行失败后的最大重试次数
distributed_ddl_allow_replicating_to_leader1是否允许在Leader节点上执行DDL

5.4 监控分布式DDL执行状态

通过 system.distributed_ddl_queue 系统表可以监控DDL任务的执行状态。

SELECT
    cluster,
    query,
    host,
    status,
    IF(error != '', error, 'OK') AS result,
    toDateTime(entry_time) AS submitted_at,
    toDateTime(finish_time) AS finished_at,
    dateDiff('second', submitted_at, finished_at) AS duration_sec
FROM system.distributed_ddl_queue
ORDER BY entry_time DESC
LIMIT 20;

字段说明

字段说明
cluster集群名称
queryDDL语句
host执行节点的主机名
status执行状态(Finished/Skipped/Unknown)
error错误信息(为空表示成功)
entry_time任务入队时间
finish_time任务完成时间

5.5 常见DDL执行问题与处理

问题1:部分节点执行失败

现象:执行 ON CLUSTER DDL后,部分节点报错。

原因

  • 节点网络不通
  • 节点磁盘空间不足
  • 节点版本不一致(DDL语法不兼容)

处理

-- 查看具体哪个节点失败
SELECT host, status, error
FROM system.distributed_ddl_queue
WHERE query LIKE '%CREATE TABLE local_user_events%'
ORDER BY entry_time DESC
LIMIT 1;

-- 在失败的节点上手动执行DDL
-- 然后清理ZK中的DDL队列(谨慎操作)
问题2:DDL任务卡住

现象:DDL执行长时间不返回。

原因

  • 某个节点宕机,ZK等待超时
  • distributed_ddl_task_timeout 设置过大

处理

-- 增大超时时间(如果集群确实很大)
SET distributed_ddl_task_timeout = 600;

-- 检查是否有节点无法访问
SELECT host_name, errors_count
FROM system.clusters
WHERE cluster = 'online_cluster';
问题3:ZK队列积压

现象:DDL执行越来越慢。

原因/clickhouse/task_queue/ddl/ 路径下积累了大量历史任务。

处理

# 清理ZK中已完成的DDL任务(保留最近N个)
# 使用zkCli.sh连接到ZooKeeper
ls /clickhouse/task_queue/ddl/online_cluster
# 手动删除已完成的任务节点

6. 分布式建表最佳实践

6.1 标准建表流程

在ClickHouse分布式集群中,正确的建表流程是:

Step 1: 在每个节点上创建本地表(ReplicatedMergeTree)
   │
   ▼
Step 2: 在每个节点上创建分布式表(Distributed引擎)
   │
   ▼
Step 3: 验证表结构一致性

为什么不能只建分布式表?

Distributed表本身不存储数据,它只是一个"视图"。如果底层没有本地表,查询Distributed表将返回空结果。

6.2 命名规范

为了区分本地表和分布式表,建议采用统一的命名规范:

表类型命名规范示例
本地表local_[table_name][table_name]_locallocal_user_events
分布式表distributed_[table_name][table_name]distributed_user_eventsuser_events
缓冲区表buffer_[table_name]buffer_user_events

推荐方案

-- 本地表:带 _local 后缀
CREATE TABLE user_events_local ON CLUSTER 'online_cluster' (...);

-- 分布式表:不带后缀(作为业务访问入口)
CREATE TABLE user_events ON CLUSTER 'online_cluster'
ENGINE = Distributed('online_cluster', 'default', 'user_events_local', cityHash64(user_id));

6.3 完整建表示例

-- ========== Step 1: 创建本地表 ==========
CREATE TABLE IF NOT EXISTS user_events_local ON CLUSTER 'online_cluster'
(
    event_date Date,
    user_id UInt64,
    event_type LowCardinality(String),
    event_value Float64,
    event_time DateTime,
    extra_info String
)
ENGINE = ReplicatedMergeTree('/clickhouse/tables/{shard}/user_events', '{replica}')
PARTITION BY toYYYYMM(event_date)
ORDER BY (user_id, event_date, event_time)
TTL event_date + INTERVAL 90 DAY
SETTINGS
    index_granularity = 8192,
    storage_configuration = 'default';

-- ========== Step 2: 创建分布式表 ==========
CREATE TABLE IF NOT EXISTS user_events ON CLUSTER 'online_cluster'
ENGINE = Distributed('online_cluster', 'default', 'user_events_local', cityHash64(user_id));

-- ========== Step 3: 验证 ==========
-- 在任意节点执行
SELECT
    table,
    engine,
    partitions,
    parts,
    rows
FROM system.tables
WHERE table IN ('user_events_local', 'user_events');

7. 实战案例:3分片2副本集群完整搭建

7.1 环境规划

节点分片副本IP地址角色
ch-node-0111192.168.1.101分片1副本1
ch-node-0212192.168.1.102分片1副本2
ch-node-0321192.168.1.103分片2副本1
ch-node-0422192.168.1.104分片2副本2
ch-node-0531192.168.1.105分片3副本1
ch-node-0632192.168.1.106分片3副本2

7.2 配置文件

config.xml(所有节点相同部分)

<clickhouse>
    <remote_servers>
        <prod_cluster>
            <shard>
                <internal_replication>true</internal_replication>
                <replica>
                    <host>192.168.1.101</host>
                    <port>9000</port>
                    <user>default</user>
                </replica>
                <replica>
                    <host>192.168.1.102</host>
                    <port>9000</port>
                    <user>default</user>
                </replica>
            </shard>
            <shard>
                <internal_replication>true</internal_replication>
                <replica>
                    <host>192.168.1.103</host>
                    <port>9000</port>
                    <user>default</user>
                </replica>
                <replica>
                    <host>192.168.1.104</host>
                    <port>9000</port>
                    <user>default</user>
                </replica>
            </shard>
            <shard>
                <internal_replication>true</internal_replication>
                <replica>
                    <host>192.168.1.105</host>
                    <port>9000</port>
                    <user>default</user>
                </replica>
                <replica>
                    <host>192.168.1.106</host>
                    <port>9000</port>
                    <user>default</user>
                </replica>
            </shard>
        </prod_cluster>
    </remote_servers>

    <!-- ZooKeeper配置 -->
    <zookeeper>
        <node index="1">
            <host>zk-01</host>
            <port>2181</port>
        </node>
        <node index="2">
            <host>zk-02</host>
            <port>2181</port>
        </node>
        <node index="3">
            <host>zk-03</host>
            <port>2181</port>
        </node>
    </zookeeper>
</clickhouse>

macros.xml(各节点不同)

<!-- ch-node-01 -->
<clickhouse>
    <macros>
        <cluster>prod_cluster</cluster>
        <shard>01</shard>
        <replica>ch-node-01</replica>
    </macros>
</clickhouse>
<!-- ch-node-02 -->
<clickhouse>
    <macros>
        <cluster>prod_cluster</cluster>
        <shard>01</shard>
        <replica>ch-node-02</replica>
    </macros>
</clickhouse>

(其他节点类似,shard和replica相应调整)

7.3 建表与验证

-- 1. 创建数据库
CREATE DATABASE IF NOT EXISTS analytics ON CLUSTER 'prod_cluster';

-- 2. 创建本地表
CREATE TABLE analytics.user_events_local ON CLUSTER 'prod_cluster'
(
    event_date Date,
    user_id UInt64,
    event_type String,
    event_value Float64,
    country LowCardinality(String),
    city String
)
ENGINE = ReplicatedMergeTree('/clickhouse/tables/{shard}/analytics/user_events', '{replica}')
PARTITION BY toYYYYMM(event_date)
ORDER BY (user_id, event_date)
SAMPLE BY user_id;

-- 3. 创建分布式表
CREATE TABLE analytics.user_events ON CLUSTER 'prod_cluster'
ENGINE = Distributed('prod_cluster', 'analytics', 'user_events_local', cityHash64(user_id));

-- 4. 写入测试数据
INSERT INTO analytics.user_events
SELECT
    toDate('2026-01-01') + rand() % 365 AS event_date,
    rand() % 1000000 AS user_id,
    ['click', 'view', 'purchase'][rand() % 3 + 1] AS event_type,
    rand() / 1000000.0 AS event_value,
    ['CN', 'US', 'JP', 'UK', 'DE'][rand() % 5 + 1] AS country,
    'city_' || toString(rand() % 100) AS city
FROM numbers(1000000);

-- 5. 验证数据分布
SELECT
    shard_num,
    count() AS rows,
    round(count() * 100.0 / sum(count()) OVER (), 2) AS pct
FROM cluster('prod_cluster', 'analytics', 'user_events_local', rand())
GROUP BY shard_num
ORDER BY shard_num;

-- 预期输出(数据大致均匀分布在3个分片):
-- shard_num | rows   | pct
--         1 | 334123 | 33.41
--         2 | 332891 | 33.29
--         3 | 333986 | 33.30

8. 总结与最佳实践

本文深入讲解了ClickHouse数据分片机制与集群配置的核心知识,以下是关键要点总结:

配置层面

  1. always set internal_replication=true:让ReplicatedMergeTree处理副本同步,避免数据不一致。
  2. 使用宏变量简化运维:通过 {shard}{replica} 等宏变量,可以用同一套DDL管理整个集群。
  3. ZooKeeper独立部署:ZK是ClickHouse分布式机制的"大脑",务必独立部署并保证其性能。

分片策略

  1. 优先使用哈希分片cityHash64(column) 是最常用且效果最好的分片方式。
  2. 避免低基数列作为分片键:如果分片键只有很少的几个不同值,会导致严重的数据倾斜。
  3. 分片数不是越多越好:每个分片都会增加查询的并行度和网络开销,需要根据数据量和查询特点合理规划。

分布式DDL

  1. 使用ON CLUSTER管理表结构:避免手动在每个节点上执行DDL,减少人为错误。
  2. 监控DDL执行状态:通过 system.distributed_ddl_queue 及时发现执行异常的节点。
  3. 控制DDL并发量:大量并发的ON CLUSTER DDL会对ZK造成压力。

建表规范

  1. 先本地表,后分布式表:确保本地表(ReplicatedMergeTree)存在后再创建Distributed表。
  2. 采用统一的命名规范:如 table_local + table(分布式视图)。
  3. 写入本地表,查询分布式表:这是ClickHouse官方推荐的最佳实践(详见下一篇文章)。

在下一篇文章中,我们将深入解析Distributed引擎的工作原理,包括分布式写入流程、分布式查询流程、GLOBAL IN/JOIN的使用场景等核心主题。


上一篇【第46篇】ClickHouse ReplicatedMergeTree原理深度解析
下一篇【第48篇】ClickHouse Distributed引擎原理_分布式读写核心流程


更多推荐