【Clickhouse从入门到精通】第47篇:ClickHouse数据分片——集群配置与分布式DDL
上一篇【第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的数据分片机制,需要掌握以下几个核心概念:
- 集群(Cluster):由多个分片(Shard)组成的逻辑集合,每个分片可以包含多个副本(Replica)。
- 分片(Shard):数据水平分割的基本单位,每个分片存储整个数据集的一个子集。
- 副本(Replica):分片的数据冗余备份,用于高可用和读负载均衡。
- 分布式表(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引擎)将数据同步到其他副本。
工作流程:
- 写入请求到达Distributed表
- Distributed表根据sharding_key计算目标分片
- 数据只写入该分片的第一个健康副本
- 该副本通过ZooKeeper协调,将数据异步同步到其他副本
优点:
- 避免重复写入,减少网络开销
- 利用ReplicatedMergeTree原生的复制机制,保证数据一致性
- 写入性能更好
缺点:
- 依赖ZooKeeper,ZK压力较大时需要关注
internal_replication = false
当设置为 false 时,ClickHouse会将数据同时写入分片内的所有副本。
工作流程:
- 写入请求到达Distributed表
- 数据被同时发送到该分片的所有副本
- 各副本独立写入,不通过ZooKeeper同步
优点:
- 不依赖ZooKeeper,适合ZK不稳定的场景
- 写入路径简单
缺点:
- 网络带宽是true模式的倍数(副本数个)
- 可能出现数据不一致(部分副本写入失败)
- 不推荐在生产环境使用
对比总结
| 特性 | internal_replication=true | internal_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.xml 或 config.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}被替换为各节点自身的分片编号(如01、02、03){replica}被替换为各节点自身的副本名称(如ch-online-01)
常用宏变量命名规范
| 宏变量名 | 含义 | 示例值 |
|---|---|---|
{cluster} | 集群名称 | online_cluster |
{shard} | 分片编号 | 01、02、03 |
{replica} | 副本标识 | ch-node-01、replica2 |
{database} | 数据库名 | analytics |
{table} | 表名 | user_events |
宏变量的高级用法
宏变量不仅可以用于ReplicatedMergeTree的ZK路径,还可以用于:
- 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';
- 字典文件的路径
<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;
字段说明:
| 字段 | 类型 | 说明 |
|---|---|---|
cluster | String | 集群名称 |
shard_num | UInt32 | 分片编号(从1开始) |
shard_weight | UInt32 | 分片权重(用于加权分片) |
replica_num | UInt32 | 副本编号(从1开始) |
host_name | String | 主机名 |
host_address | String | IP地址 |
port | UInt16 | 端口号 |
user | String | 用户名 |
is_local | UInt8 | 是否为本地节点 |
errors_count | UInt32 | 连接错误次数 |
estimated_recovery_time | UInt32 | 预计恢复时间(秒) |
示例输出:
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) | ⭐⭐⭐ | ✅ | ✅ | ⭐ | ⭐⭐⭐ |
| 范围分片 | ⭐⭐ | ✅✅ | 取决于范围 | ⭐⭐⭐ | ⭐⭐ |
| 一致性哈希 | ⭐⭐⭐⭐ | ✅ | ✅ | ⭐⭐⭐⭐ | ⭐⭐⭐⭐ |
数据倾斜风险分析:
- 随机分片:理论上无数据倾斜,但实际中由于随机数分布特性,小数据量时可能有偏差。
- 哈希分片:如果分片键的基数低(如只有几个不同值),会导致严重的数据倾斜。
- 取模分片:当分片数变化时,大部分数据需要重新分片。
- 范围分片:数据分布最不均匀,容易出现热点分片。
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_timeout | 180 | DDL任务等待所有节点完成的超时时间 |
distributed_ddl_task_max_tries | 3 | 执行失败后的最大重试次数 |
distributed_ddl_allow_replicating_to_leader | 1 | 是否允许在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 | 集群名称 |
query | DDL语句 |
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]_local | local_user_events |
| 分布式表 | distributed_[table_name] 或 [table_name] | distributed_user_events 或 user_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-01 | 1 | 1 | 192.168.1.101 | 分片1副本1 |
| ch-node-02 | 1 | 2 | 192.168.1.102 | 分片1副本2 |
| ch-node-03 | 2 | 1 | 192.168.1.103 | 分片2副本1 |
| ch-node-04 | 2 | 2 | 192.168.1.104 | 分片2副本2 |
| ch-node-05 | 3 | 1 | 192.168.1.105 | 分片3副本1 |
| ch-node-06 | 3 | 2 | 192.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数据分片机制与集群配置的核心知识,以下是关键要点总结:
配置层面
- always set
internal_replication=true:让ReplicatedMergeTree处理副本同步,避免数据不一致。 - 使用宏变量简化运维:通过
{shard}、{replica}等宏变量,可以用同一套DDL管理整个集群。 - ZooKeeper独立部署:ZK是ClickHouse分布式机制的"大脑",务必独立部署并保证其性能。
分片策略
- 优先使用哈希分片:
cityHash64(column)是最常用且效果最好的分片方式。 - 避免低基数列作为分片键:如果分片键只有很少的几个不同值,会导致严重的数据倾斜。
- 分片数不是越多越好:每个分片都会增加查询的并行度和网络开销,需要根据数据量和查询特点合理规划。
分布式DDL
- 使用ON CLUSTER管理表结构:避免手动在每个节点上执行DDL,减少人为错误。
- 监控DDL执行状态:通过
system.distributed_ddl_queue及时发现执行异常的节点。 - 控制DDL并发量:大量并发的ON CLUSTER DDL会对ZK造成压力。
建表规范
- 先本地表,后分布式表:确保本地表(ReplicatedMergeTree)存在后再创建Distributed表。
- 采用统一的命名规范:如
table_local+table(分布式视图)。 - 写入本地表,查询分布式表:这是ClickHouse官方推荐的最佳实践(详见下一篇文章)。
在下一篇文章中,我们将深入解析Distributed引擎的工作原理,包括分布式写入流程、分布式查询流程、GLOBAL IN/JOIN的使用场景等核心主题。
上一篇【第46篇】ClickHouse ReplicatedMergeTree原理深度解析
下一篇【第48篇】ClickHouse Distributed引擎原理_分布式读写核心流程
更多推荐
所有评论(0)