ClickHouse Keeper 为数据复制和分布式 DDL查询的执行提供协调系统。ClickHouse Keeper 与 ZooKeeper 兼容。
实现细节
ZooKeeper 是最早广为人知的开源协调系统之一。它使用 Java 实现,并拥有简单而强大的数据模型。ZooKeeper 的协调算法 ZooKeeper Atomic Broadcast (ZAB) 无法为读操作提供线性一致性保证,因为每个 ZooKeeper 节点都在本地处理读请求。与 ZooKeeper 不同,ClickHouse Keeper 使用 C++ 编写,并采用了 RAFT 算法的实现。该算法可为读写同时提供线性一致性,并且在多种语言中都有开源实现。
默认情况下,ClickHouse Keeper 提供与 ZooKeeper 相同的保证:写入具备线性一致性,而读取不具备线性一致性。它具有兼容的客户端/服务器协议,因此任何标准的 ZooKeeper 客户端都可以与 ClickHouse Keeper 交互。快照和日志的格式与 ZooKeeper 不兼容,但 clickhouse-keeper-converter 工具可以将 ZooKeeper 数据转换为 ClickHouse Keeper 快照。ClickHouse Keeper 的服务器间协议也与 ZooKeeper 不兼容,因此无法构建混合的 ZooKeeper / ClickHouse Keeper 集群。
ClickHouse Keeper 对访问控制列表 (ACL) 的支持方式与 ZooKeeper 相同。ClickHouse Keeper 支持同样的一组权限,并提供完全相同的内置方案:world、auth 和 digest。digest 身份验证方案使用 username:password 这一组凭据,其中密码采用 Base64 编码。
配置
ClickHouse Keeper 既可作为 ZooKeeper 的独立替代方案使用,也可作为 ClickHouse server 的内部组件。在这两种情况下,配置使用的 .xml 文件几乎相同。
Keeper 配置设置
主要的 ClickHouse Keeper 配置标签为 <keeper_server>,包含以下参数:
| Parameter | Description | Default |
|---|---|---|
tcp_port |
客户端连接使用的端口。 | 2181 |
tcp_port_secure |
用于客户端与 keeper-server 之间建立 SSL 连接的安全端口。 | - |
server_id |
唯一的服务器 ID,ClickHouse Keeper 集群中的每个参与者都必须具有唯一编号 (1、2、3……) 。 | - |
log_storage_path |
协调日志的存储路径。与 ZooKeeper 类似,最佳做法是将日志存储在负载较低的节点上。 | - |
snapshot_storage_path |
协调快照的存储路径。 | - |
enable_reconfiguration |
通过 reconfig 启用集群动态重配置。 |
False |
max_memory_usage_soft_limit |
keeper 最大内存使用量的软限制,单位为字节。 | max_memory_usage_soft_limit_ratio * physical_memory_amount |
max_memory_usage_soft_limit_ratio |
如果未设置 max_memory_usage_soft_limit 或将其设置为零,则使用该值来定义默认软限制。 |
0.9 |
cgroups_memory_observer_wait_time |
如果未设置 max_memory_usage_soft_limit 或其值为 0,则使用该时间间隔来监测物理内存总量。一旦内存总量发生变化,就会根据 max_memory_usage_soft_limit_ratio 重新计算 Keeper 的内存软限制。 |
15 |
http_control |
HTTP 控制 接口的配置。 | - |
digest_enabled |
启用实时数据一致性检查。 | True |
create_snapshot_on_exit |
在关闭时创建快照。 | - |
hostname_checks_enabled |
为集群配置启用主机名合理性检查 (例如,在远程端点中使用 localhost 时) 。 | True |
four_letter_word_white_list |
4lw 命令白名单。 | conf, cons, crst, envi, ruok, srst, srvr, stat, wchs, dirs, mntr, isro, rcvr, apiv, csnp, lgif, rqld, ydld |
enable_ipv6 |
启用 IPv6。 | True |
其他常用参数继承自 ClickHouse server 配置 (listen_host、logger 等) 。
内部协调设置
内部协调设置位于 <keeper_server>.<coordination_settings> 部分,包含以下参数:
| 参数 | 说明 | 默认值 |
|---|---|---|
operation_timeout_ms |
单个客户端操作的超时时间 (毫秒) | 10000 |
min_session_timeout_ms |
客户端会话的最小超时时间 (毫秒) | 10000 |
session_timeout_ms |
客户端会话的最大超时时间 (毫秒) | 100000 |
dead_session_check_period_ms |
ClickHouse Keeper 检查失效会话并将其移除的频率 (毫秒) | 500 |
heart_beat_interval_ms |
ClickHouse Keeper 的 leader 向跟随者发送心跳的频率 (毫秒) | 500 |
election_timeout_lower_bound_ms |
如果跟随者在此时间间隔内未收到 leader 的心跳,则可以发起 leader 选举。该值必须小于或等于 election_timeout_upper_bound_ms。理想情况下,两者不应相等。 |
1000 |
election_timeout_upper_bound_ms |
如果跟随者在此时间间隔内未收到 leader 的心跳,则必须发起 leader 选举。 | 2000 |
rotate_log_storage_interval |
单个文件中存储的日志记录数。 | 100000 |
reserved_log_items |
执行合并整理前保留的协调日志记录数。 | 100000 |
snapshot_distance |
ClickHouse Keeper 创建新快照的频率 (按日志记录数计算) 。 | 100000 |
snapshots_to_keep |
保留的快照数量。 | 3 |
stale_log_gap |
当 leader 认为跟随者已过时,并向其发送快照而不是日志时使用的阈值。 | 10000 |
fresh_log_gap |
节点变为最新状态时的阈值。 | 200 |
max_requests_batch_size |
请求发送到 RAFT 前,按请求数计算的批次最大大小。 | 100 |
force_sync |
每次写入协调日志时都调用 fsync。 |
true |
quorum_reads |
通过整个 RAFT 共识流程执行读取请求,将其按写入方式处理,速度相近。 | false |
raft_logs_level |
协调相关的文本日志级别 (trace、debug 等) 。 | system default |
auto_forwarding |
允许将写入请求从跟随者转发给 leader。 | true |
shutdown_timeout |
等待内部连接完成并关闭的时间 (毫秒) 。 | 5000 |
startup_timeout |
如果服务器未在指定超时时间内连接到其他 quorum 参与者,则会终止 (毫秒) 。 | 30000 |
async_replication |
启用异步复制。在提升性能的同时保留所有读写保证。此设置默认禁用,以避免破坏向后兼容性。 | false |
latest_logs_cache_size_threshold |
最新日志条目的内存缓存最大总大小。 | 1GiB |
commit_logs_cache_size_threshold |
已弃用。如果未设置 log_readahead_commit_window_bytes,则使用此设置。 |
- |
commit_logs_cache_entry_count_threshold |
已弃用,无任何作用。请改用 log_readahead_commit_window_bytes。 |
- |
log_readahead_enabled |
为更新日志追赶读取 (跟随者复制) 启用按 peer 解码的预读。 | false |
log_readahead_window_bytes |
每个 peer reader 缓冲的已解码条目的最大字节数。应至少与典型的 append-entries 批次一样大。 | 64MiB |
log_readahead_max_peer_readers |
并发的每 peer 预读 reader 的最大数量。 | 8 |
log_readahead_eviction_timeout_ms |
非活动的每 peer 预读 reader 或 commit reader 被逐出前的空闲超时时间 (后续未命中会重新创建它) 。 | 30000 |
log_readahead_pool_threads |
专用预读线程池中的线程数。0 会从 log_readahead_max_peer_readers 推导该值。 |
0 |
log_readahead_serve_wait_timeout_ms |
回退到直接读取前,等待后台预读填充的最长时间。 | 200 |
log_readahead_chunk_size |
预读填充任务中每个块解码的日志条目数。 | 16 |
log_readahead_commit_window_bytes |
提交线程前预先缓冲的已解码日志条目的最大总大小。0 会禁用提交预读 (提交操作会逐个从磁盘读取条目) 。 |
500MiB |
log_startup_read_max_streams |
Keeper 启动期间并发读取的更新日志文件最大数量。0 = 自动使用 CPU 核心数。1 = 使用串行 (并行化前) 的启动读取。有效并行度受需读取的更新日志文件数量限制;对于寻道受限的存储 (HDD、IOPS 受限的卷) ,请考虑降低此值。 |
0 |
log_startup_read_buffer_size |
Keeper 启动时读取更新日志所使用的每流读取缓冲区大小 (字节) 。必须大于 0。缓冲区还会被限制为不超过文件大小。 |
8MiB |
disk_move_retries_wait_ms |
文件在磁盘之间移动时发生失败后,两次重试之间的等待时间。 | 1000 |
disk_move_retries_during_init |
初始化期间,文件在磁盘之间移动时发生失败后的重试次数。 | 100 |
experimental_use_rocksdb |
使用 RocksDB 作为后端存储。 | 0 |
quorum 的配置位于 <keeper_server>.<raft_configuration> 部分,其中包含服务器描述。
整个 quorum 唯一的参数是 secure,用于启用 quorum 参与者之间通信的加密连接。如果节点之间的内部通信需要 SSL 连接,可将该参数设为 true;否则可不指定。
每个 <server> 的主要参数为:
id— quorum 中的服务器标识符。hostname— 该服务器所在主机的主机名。port— 该服务器监听连接的端口。can_become_leader— 设为false可将服务器配置为learner。如果省略,则默认为true。
带有 test_keeper_ 前缀的集成测试中提供了三节点 quorum 的配置示例。下面是服务器 #1 的示例配置:
<keeper_server>
<tcp_port>2181</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>trace</raft_logs_level>
</coordination_settings>
<raft_configuration>
<server>
<id>1</id>
<hostname>zoo1</hostname>
<port>9234</port>
</server>
<server>
<id>2</id>
<hostname>zoo2</hostname>
<port>9234</port>
</server>
<server>
<id>3</id>
<hostname>zoo3</hostname>
<port>9234</port>
</server>
</raft_configuration>
</keeper_server>如何运行
ClickHouse Keeper 已集成在 ClickHouse server 软件包中,只需将 <keeper_server> 配置添加到 /etc/your_path_to_config/clickhouse-server/config.xml,然后像平常一样启动 ClickHouse server 即可。如果你想运行独立运行的 ClickHouse Keeper,也可以用类似的方式启动:
clickhouse-keeper --config /etc/your_path_to_config/config.xml如果没有该符号链接 (clickhouse-keeper) ,您可以创建一个,或将 keeper 指定为 clickhouse 的参数:
clickhouse keeper --config /etc/your_path_to_config/config.xml四字母命令
ClickHouse Keeper 也提供 4lw 命令,其用法与 ZooKeeper 基本相同。每个命令都由四个字母组成,例如 mntr、stat 等。还有一些比较常用的命令:stat 提供有关服务器和已连接客户端的常规信息,srvr 提供服务器的扩展信息,cons 提供连接的扩展信息。
4lw 命令有一个白名单配置项 four_letter_word_white_list,其默认值为 conf,cons,crst,envi,ruok,srst,srvr,stat,wchs,dirs,mntr,isro,rcvr,apiv,csnp,lgif,rqld,ydld。
你可以通过 telnet 或 nc 在客户端端口向 ClickHouse Keeper 发送这些命令。
echo mntr | nc localhost 9181以下是 4lw 命令的详细说明:
ruok:测试服务器是否正在运行且未处于 error 状态。如果服务器正在运行,会返回imok;否则将完全没有响应。返回imok并不一定表示服务器已加入 quorum,它只说明服务器进程处于活动状态,并已绑定到指定的客户端端口。要查看与 quorum 相关的状态以及客户端连接信息,请使用stat。
imokmntr:输出一组可用于监控集群健康状态的变量。
zk_version v21.11.1.1-prestable-7a4a0b0edef0ad6e0aa662cd3b90c3f4acf796e7
zk_avg_latency 0
zk_max_latency 0
zk_min_latency 0
zk_packets_received 68
zk_packets_sent 68
zk_num_alive_connections 1
zk_outstanding_requests 0
zk_server_state leader
zk_znode_count 4
zk_watch_count 1
zk_ephemerals_count 0
zk_approximate_data_size 723
zk_open_file_descriptor_count 310
zk_max_file_descriptor_count 10240
zk_leader_uptime 1234
zk_sum_leader_unavailable_time 1234
zk_cnt_leader_unavailable_time 1
zk_sum_election_time 1234
zk_cnt_election_time 1
zk_followers 0
zk_synced_followers 0zk_sum_leader_unavailable_time、zk_cnt_leader_unavailable_time、zk_sum_election_time 和 zk_cnt_election_time 是仅针对 leader 的累计指标。它们是每个服务器的观测值,而非集群级可用性度量:当该服务器不再观测到存活的 leader 时,便开始计时。在网络分区期间,其中可能包括旧 leader 在另一分区中仍然存活的时间。选举指标仅记录在此类本地观测到的无 leader 窗口之后成功完成的选举;未出现采样到的无 leader 状态的领导权转移不会被计入,这是有意的设计。
Keeper 每隔 heart_beat_interval_ms 对其本地 NuRaft leader 状态采样一次,但采样频率不会高于每 100 毫秒一次。每个边界与本地状态转换之间最多可能相差一个有效时间间隔,短于该时间间隔的窗口可能会被遗漏。Keeper 会在 NuRaft 的 BecomeLeader 时记录本地选举完成。srst 会重置这四个值。
每种最近一次本地观测到的窗口耗时也会以毫秒为单位作为 KeeperLastLeaderElectionTime 和 KeeperLastLeaderUnavailableTime 导出到 system.asynchronous_metrics。KeeperLastLeaderElectionTime 遵循与 zk_sum_election_time 相同的无 leader 窗口定义。两者也都仅针对 leader:非活动 leader 的节点,或尚未完成窗口的节点,都会报告 0。srst 会将它们与累计计数器一起重置。
srvr:列出服务器的完整信息。
ClickHouse Keeper version: v21.11.1.1-prestable-7a4a0b0edef0ad6e0aa662cd3b90c3f4acf796e7
Latency min/avg/max: 0/0/0
Received: 2
Sent : 2
Connections: 1
Outstanding: 0
Zxid: 34
Mode: leader
Node count: 4stat:列出服务器和已连接客户端的简要信息。
ClickHouse Keeper version: v21.11.1.1-prestable-7a4a0b0edef0ad6e0aa662cd3b90c3f4acf796e7
Clients:
192.168.1.1:52852(recved=0,sent=0)
192.168.1.1:52042(recved=24,sent=48)
Latency min/avg/max: 0/0/0
Received: 4
Sent : 4
Connections: 1
Outstanding: 0
Zxid: 36
Mode: leader
Node count: 4srst:重置服务器统计信息。该命令会影响srvr、mntr和stat的结果。
Server stats reset.conf:输出服务配置的详细信息。
server_id=1
tcp_port=2181
four_letter_word_white_list=*
log_storage_path=./coordination/logs
snapshot_storage_path=./coordination/snapshots
max_requests_batch_size=100
session_timeout_ms=30000
operation_timeout_ms=10000
dead_session_check_period_ms=500
heart_beat_interval_ms=500
election_timeout_lower_bound_ms=1000
election_timeout_upper_bound_ms=2000
reserved_log_items=1000000000000000
snapshot_distance=10000
auto_forwarding=true
shutdown_timeout=5000
startup_timeout=240000
raft_logs_level=information
snapshots_to_keep=3
rotate_log_storage_interval=100000
stale_log_gap=10000
fresh_log_gap=200
max_requests_batch_size=100
quorum_reads=false
force_sync=false
compress_logs=true
compress_snapshots_with_zstd_format=true
configuration_change_tries_count=20cons:列出连接到此服务器的所有客户端的完整连接/会话信息。包括收发数据包数量、会话 ID、操作延迟、最近执行的操作等信息……
192.168.1.1:52163(recved=0,sent=0,sid=0xffffffffffffffff,lop=NA,est=1636454787393,to=30000,lzxid=0xffffffffffffffff,lresp=0,llat=0,minlat=0,avglat=0,maxlat=0)
192.168.1.1:52042(recved=9,sent=18,sid=0x0000000000000001,lop=List,est=1636454739887,to=30000,lcxid=0x0000000000000005,lzxid=0x0000000000000005,lresp=1636454739892,llat=0,minlat=0,avglat=0,maxlat=0)crst:重置所有连接的连接/会话统计信息。
Connection stats reset.envi: 输出服务环境的详细信息
Environment:
clickhouse.keeper.version=v21.11.1.1-prestable-7a4a0b0edef0ad6e0aa662cd3b90c3f4acf796e7
host.name=ZBMAC-C02D4054M.local
os.name=Darwin
os.arch=x86_64
os.version=19.6.0
cpu.count=12
user.name=root
user.home=/Users/JackyWoo/
user.dir=/Users/JackyWoo/project/jd/clickhouse/cmake-build-debug/programs/
user.tmp=/var/folders/b4/smbq5mfj7578f2jzwn602tt40000gn/T/dirs:显示快照和日志文件的总大小 (单位为字节)
snapshot_dir_size: 0
log_dir_size: 3875isro:用于检查 server 是否以只读模式运行。处于只读模式时,server 会返回ro;否则返回rw。
rwwchs:列出服务器上 watches 的简要信息。
1 connections watching 1 paths
Total watches:1wchc:按会话列出服务器上的 watch 详细信息。该命令会输出会话 (连接) 及其关联的 watch (路径) 列表。请注意,具体开销取决于 watch 的数量,此操作可能代价较高 (会影响服务器性能) ,请谨慎使用。
0x0000000000000001
/clickhouse/task_queue/ddlwchp:按路径列出服务器上的 watches 详细信息。它会输出路径 (znodes) 及其关联会话的列表。请注意,watches 数量较多时,此操作的开销可能很大 (即会影响服务器性能) ,请谨慎使用。
/clickhouse/task_queue/ddl
0x0000000000000001dump:列出未完成的会话和临时节点。此命令仅在 leader 上有效。
Sessions dump (2):
0x0000000000000001
0x0000000000000002
Sessions with Ephemerals (1):
0x0000000000000001
/clickhouse/task_queue/ddlcsnp:调度一个快照创建任务。成功时返回已调度快照的最后一个已提交日志索引,失败时返回Failed to schedule snapshot creation task.。lgif命令可帮助你判断快照是否已完成。
100lgif:Keeper 日志信息。first_log_idx:当前节点在日志存储中的第一个日志索引;first_log_term:当前节点的第一个日志任期;last_log_idx:当前节点在日志存储中的最后一个日志索引;last_log_term:当前节点的最后一个日志任期;last_committed_log_idx:当前节点在状态机中的最后一个已提交日志索引;leader_committed_log_idx:当前节点视角下 leader 的已提交日志索引;target_committed_log_idx:应提交到的目标日志索引;last_snapshot_idx:最近一次快照中最大的已提交日志索引。
first_log_idx 1
first_log_term 1
last_log_idx 101
last_log_term 1
last_committed_log_idx 100
leader_committed_log_idx 101
target_committed_log_idx 101
last_snapshot_idx 50rqld:请求成为新的 leader。请求已发送时返回Sent leadership request to leader.,未发送时返回Failed to send leadership request to leader.。如果该节点已经是 leader,结果与请求已发送时相同。
Sent leadership request to leader.ftfl:列出所有功能开关及其在该 Keeper 实例中是否已启用。
filtered_list 1
multi_read 1
check_not_exists 0ydld:请求让出 leader 身份并转为 follower。如果接收该请求的服务器是 leader,它会先暂停写操作,等待继任者 (当前 leader 不会成为继任者) 追赶至最新日志,然后辞去 leader 身份。继任者将自动选出。请求已发送时返回Sent yield leadership request to leader.,未发送时返回Failed to send yield leadership request to leader.。如果该节点已是 follower,结果与发送该请求时相同。
Sent yield leadership request to leader.pfev:返回所有已收集事件的值。对于每个事件,返回事件名称、事件值和事件描述。
FileOpen 62 Number of files opened.
Seek 4 Number of times the 'lseek' function was called.
ReadBufferFromFileDescriptorRead 126 Number of reads (read/pread) from a file descriptor. Does not include sockets.
ReadBufferFromFileDescriptorReadFailed 0 Number of times the read (read/pread) from a file descriptor have failed.
ReadBufferFromFileDescriptorReadBytes 178846 Number of bytes read from file descriptors. If the file is compressed, this will show the compressed data size.
WriteBufferFromFileDescriptorWrite 7 Number of writes (write/pwrite) to a file descriptor. Does not include sockets.
WriteBufferFromFileDescriptorWriteFailed 0 Number of times the write (write/pwrite) to a file descriptor have failed.
WriteBufferFromFileDescriptorWriteBytes 153 Number of bytes written to file descriptors. If the file is compressed, this will show compressed data size.
FileSync 2 Number of times the F_FULLFSYNC/fsync/fdatasync function was called for files.
DirectorySync 0 Number of times the F_FULLFSYNC/fsync/fdatasync function was called for directories.
FileSyncElapsedMicroseconds 12756 Total time spent waiting for F_FULLFSYNC/fsync/fdatasync syscall for files.
DirectorySyncElapsedMicroseconds 0 Total time spent waiting for F_FULLFSYNC/fsync/fdatasync syscall for directories.
ReadCompressedBytes 0 Number of bytes (the number of bytes before decompression) read from compressed sources (files, network).
CompressedReadBufferBlocks 0 Number of compressed blocks (the blocks of data that are compressed independent of each other) read from compressed sources (files, network).
CompressedReadBufferBytes 0 Number of uncompressed bytes (the number of bytes after decompression) read from compressed sources (files, network).
AIOWrite 0 Number of writes with Linux or FreeBSD AIO interface
AIOWriteBytes 0 Number of bytes written with Linux or FreeBSD AIO interface
...HTTP 控制
ClickHouse Keeper 提供了 HTTP 接口,用于检查副本是否已准备好接收流量。在 Kubernetes 等云环境中可使用该接口。
启用 /ready 端点的配置示例:
<clickhouse>
<keeper_server>
<http_control>
<port>9182</port>
<readiness>
<endpoint>/ready</endpoint>
</readiness>
</http_control>
</keeper_server>
</clickhouse>功能开关
Keeper 与 ZooKeeper 及其客户端完全兼容,但也引入了一些 ClickHouse 客户端可用的特有功能和请求类型。
由于这些功能可能带来不兼容的变更,因此其中大多数默认处于禁用状态,可通过 keeper_server.feature_flags 配置启用。
也可以显式禁用所有功能。
如果您想为 Keeper 集群启用某个新功能,建议先将集群中的所有 Keeper 实例升级到支持该功能的版本,再启用该功能本身。
下面是一个功能开关配置示例,其中禁用了 multi_read,并启用了 check_not_exists:
<clickhouse>
<keeper_server>
<feature_flags>
<multi_read>0</multi_read>
<check_not_exists>1</check_not_exists>
</feature_flags>
</keeper_server>
</clickhouse>以下功能可用:
| 功能 | 描述 | 默认值 |
|---|---|---|
multi_read |
支持 read multi 请求 | 1 |
filtered_list |
支持 list 请求,可按节点类型 (ephemeral 或 persistent) 过滤结果 | 1 |
check_not_exists |
支持 CheckNotExists 请求,用于断言节点不存在 |
1 |
create_if_not_exists |
支持 CreateIfNotExists 请求;如果节点不存在,则尝试创建该节点。如果节点已存在,则不会进行任何更改,并返回 ZOK |
1 |
remove_recursive |
支持 RemoveRecursive 请求,用于删除节点及其整个子树 |
1 |
从 ZooKeeper 迁移
无法从 ZooKeeper 无缝迁移到 ClickHouse Keeper。您必须先停止 ZooKeeper 集群,转换数据,然后再启动 ClickHouse Keeper。clickhouse-keeper-converter 工具可将 ZooKeeper 的日志和快照转换为 ClickHouse Keeper 快照。该工具要求 ZooKeeper 版本为 3.4 或更高。
迁移前准备
迁移需要停止数据摄取。请在开始前安排好维护窗口。
在停止 ZooKeeper 之前,先停止会修改协调元数据的 ClickHouse 后台任务。例如:
SYSTEM STOP MERGES;在迁移前记录用于比对的指标,以便在迁移后验证一致性。
迁移步骤
-
停止向所有 ClickHouse 节点摄取数据。
-
停止所有 ClickHouse 节点上的后台任务 (参见上文) 。
-
停止所有 ZooKeeper 节点。
-
可选但建议执行:找到 ZooKeeper 的 leader 节点,将其重新启动后再停止。这样会强制 ZooKeeper 在转换前先将一致性快照写入磁盘。
-
在 leader 节点上运行
clickhouse-keeper-converter。如果你安装了完整的 ClickHouse 可执行文件,请改用keeper-converter子命令 (clickhouse keeper-converter) 。如果两者都不可用,请下载该可执行文件。
clickhouse-keeper-converter \
--zookeeper-logs-dir /var/lib/zookeeper/version-2 \
--zookeeper-snapshots-dir /var/lib/zookeeper/version-2 \
--output-dir /path/to/clickhouse/keeper/snapshots-
将快照复制到所有 ClickHouse Keeper 节点。在任何节点启动之前,每个节点上都必须已有该快照——如果某个节点在没有快照的情况下启动,可能会以空状态将自己选举为 leader。
-
更新 ClickHouse 配置,使其指向新的 Keeper 集群。
-
在所有节点上启动 ClickHouse Keeper,然后重启 ClickHouse。
-
将各项指标与迁移前的基线进行比较,以验证一致性。
-
恢复后台任务并重启数据摄取。
整合多个 ZooKeeper 集群
如果您运行了多个 ZooKeeper 集群 (例如每个分片组一个) ,可以将它们整合到一个 ClickHouse Keeper 集群中。官方 clickhouse-keeper-converter 工具仅支持一对一转换 (一个 ZooKeeper 集群对应一个 Keeper 快照) ,因此若要进行整合,需要修改转换器的源代码以合并多个快照:
- 在每个 ZooKeeper 集群上分别运行
clickhouse-keeper-converter,并将各自的输出写入不同的目录。 - 按顺序对快照文件进行反序列化。合并时,重新计算
numChildren值,以避免不同源集群的命名空间之间发生节点 ID 冲突。 - 将合并后的输出写入目标 ClickHouse Keeper 快照目录。
处理加密和 ACL
ClickHouse Keeper 支持与 ZooKeeper 相同的 ACL 方案 (world、auth、digest) 。在转换过程中如何处理 ACL,取决于你的 ZooKeeper 配置:
- 完全加密或完全未加密:可直接转换。转换器会保留现有的 ACL 信息。
- 部分加密:转换前,先为超级管理员账户授予权限,并使用
setAcl -R清除受影响 path 上的 ACL。转换完成后,如有需要,再在 ClickHouse Keeper 中重新启用加密。
验证迁移
启动 ClickHouse Keeper 并重启 ClickHouse 后,请将关键指标与迁移前的基线进行对比,以确认迁移已成功。
在整合多个 ZooKeeper 集群时,需要区分以下两类:
- 公共路径:多个源集群中都存在且数据相同的路径——这些路径应在合并后的输出中去重。
- 差异化路径:仅存在于特定集群下的路径 (例如每个分片组在
/clickhouse/tables下的路径) ——这些路径必须从正确的源中保留下来。
避免直接遍历大型 ZooKeeper 树进行比较。应在转换过程中将所有已转换的路径输出到文件中。
迁移后调优
迁移后,若集群规模更大或吞吐量更高,可考虑调整以下设置:
| 设置 | 默认值 | 推荐值 | 说明 |
|---|---|---|---|
max_requests_batch_size |
100 |
10000 |
对于分片数较多或分片较多的集群,建议增大该值 |
force_sync |
true |
false |
异步写入日志可提高吞吐量 |
compress_logs |
false |
true |
压缩 Raft 日志文件以减少磁盘 I/O |
compress_snapshots_with_zstd_format |
true |
— | 默认已启用;使用 zstd 压缩快照 |
这些设置需在 Keeper 配置 的 coordination_settings 下进行配置。
丢失 quorum 后的恢复
由于 ClickHouse Keeper 使用 Raft,因此它可根据集群规模容忍一定数量的节点崩溃。 例如,对于一个 3 节点集群,如果只有 1 个节点崩溃,它仍能继续正常工作。
集群配置可以动态调整,但也有一些限制。重新配置同样依赖于 Raft, 因此,要向集群中添加或移除节点,就必须具备 quorum。如果集群同时丢失了过多节点,且无法 再次将其启动,Raft 就会停止工作,并且不再允许你通过常规方式重新配置集群。
不过,ClickHouse Keeper 提供了一种恢复模式,允许你仅用 1 个节点强制重新配置集群。 只有在你无法再次启动这些节点,或者无法在相同端点上启动新的实例时,才应将其作为最后手段。
继续之前,请务必注意以下事项:
- 确保故障节点无法再次连接到集群。
- 在步骤明确说明之前,不要启动任何新节点。
确认以上事项无误后,你需要执行以下操作:
- 选择一个 Keeper 节点作为新的 leader。请注意,该节点的数据将用于整个集群,因此建议选择状态最新的节点。
- 在执行任何其他操作之前,先备份所选节点的
log_storage_path和snapshot_storage_path文件夹。 - 在你打算使用的所有节点上重新配置集群。
- 向你选定的节点发送四字命令
rcvr,这会将该节点切换到恢复模式;或者停止所选节点上的 Keeper 实例,然后使用--force-recovery参数重新启动它。 - 逐个启动新节点上的 Keeper 实例,并确保在启动下一个节点之前,
mntr对zk_server_state返回follower。 - 在恢复模式下,leader 节点会对
mntr命令返回错误消息,直到它与新节点达成 quorum;同时会拒绝来自客户端和跟随者的任何请求。 - 达成 quorum 后,leader 节点将恢复正常运行模式,并开始接受所有请求。可使用
mntr通过zk_server_state返回leader来验证 Raft 已恢复正常。
在 Keeper 中使用磁盘
Keeper 支持使用部分外部磁盘来存储快照、日志文件和状态文件。
支持的磁盘类型包括:
- s3_plain
- s3
- local
下面是配置中包含的磁盘定义示例。
<clickhouse>
<storage_configuration>
<disks>
<log_local>
<type>local</type>
<path>/var/lib/clickhouse/coordination/logs/</path>
</log_local>
<log_s3_plain>
<type>s3_plain</type>
<endpoint>https://some_s3_endpoint/logs/</endpoint>
<access_key_id>ACCESS_KEY</access_key_id>
<secret_access_key>SECRET_KEY</secret_access_key>
</log_s3_plain>
<snapshot_local>
<type>local</type>
<path>/var/lib/clickhouse/coordination/snapshots/</path>
</snapshot_local>
<snapshot_s3_plain>
<type>s3_plain</type>
<endpoint>https://some_s3_endpoint/snapshots/</endpoint>
<access_key_id>ACCESS_KEY</access_key_id>
<secret_access_key>SECRET_KEY</secret_access_key>
</snapshot_s3_plain>
<state_s3_plain>
<type>s3_plain</type>
<endpoint>https://some_s3_endpoint/state/</endpoint>
<access_key_id>ACCESS_KEY</access_key_id>
<secret_access_key>SECRET_KEY</secret_access_key>
</state_s3_plain>
</disks>
</storage_configuration>
</clickhouse>要将某个磁盘用于存储日志,应将 keeper_server.log_storage_disk 配置设为该磁盘的名称。
要将某个磁盘用于存储快照,应将 keeper_server.snapshot_storage_disk 配置设为该磁盘的名称。
此外,keeper_server.latest_log_storage_disk 可用于存储最新日志,keeper_server.latest_snapshot_storage_disk 可用于存储最新快照。
在这种情况下,创建新的日志或快照时,Keeper 会自动将文件移动到正确的磁盘。
要将某个磁盘用于存储状态文件,应将 keeper_server.state_storage_disk 配置设为该磁盘的名称。
在不同磁盘之间移动文件是安全的,即使 Keeper 在传输过程中停止,也不会有数据丢失的风险。 在文件完全移动到新磁盘之前,它不会从旧磁盘中删除。
当 Keeper 的 keeper_server.coordination_settings.force_sync 设置为 true (默认值即为 true) 时,对于所有类型的磁盘,都无法满足某些保证。
目前,只有 local 类型的磁盘支持持久化同步。
如果使用 force_sync,且未使用 latest_log_storage_disk,则 log_storage_disk 应为 local 磁盘。
如果使用了 latest_log_storage_disk,则它必须始终是 local 磁盘。
如果禁用 force_sync,则在任何配置下都可以使用所有类型的磁盘。
Keeper 实例的一种可能存储配置如下所示:
<clickhouse>
<keeper_server>
<log_storage_disk>log_s3_plain</log_storage_disk>
<latest_log_storage_disk>log_local</latest_log_storage_disk>
<snapshot_storage_disk>snapshot_s3_plain</snapshot_storage_disk>
<latest_snapshot_storage_disk>snapshot_local</latest_snapshot_storage_disk>
</keeper_server>
</clickhouse>此实例会将除最新日志外的所有日志存储在 log_s3_plain 磁盘上,而最新日志则存储在 log_local 磁盘上。
快照同理:除最新快照外的所有快照都将存储在 snapshot_s3_plain 上,而最新快照则存储在 snapshot_local 磁盘上。
更改磁盘配置
如果定义了分层磁盘配置 (为最新文件使用单独的磁盘) ,Keeper 会在启动时尝试自动将文件移动到对应的磁盘。 仍然提供与之前相同的保障:在文件完全迁移到新磁盘之前,不会从旧磁盘中删除,因此可以安全地 多次重启。
如果需要将文件迁移到全新的磁盘 (或者从双磁盘配置迁移到单磁盘配置) ,可以为 keeper_server.old_snapshot_storage_disk 和 keeper_server.old_log_storage_disk 定义多个值。
以下配置展示了如何从之前的双磁盘配置迁移到全新的单磁盘配置:
<clickhouse>
<keeper_server>
<old_log_storage_disk>log_local</old_log_storage_disk>
<old_log_storage_disk>log_s3_plain</old_log_storage_disk>
<log_storage_disk>log_local2</log_storage_disk>
<old_snapshot_storage_disk>snapshot_s3_plain</old_snapshot_storage_disk>
<old_snapshot_storage_disk>snapshot_local</old_snapshot_storage_disk>
<snapshot_storage_disk>snapshot_local2</snapshot_storage_disk>
</keeper_server>
</clickhouse>启动时,所有日志文件都会从 log_local 和 log_s3_plain 移至 log_local2 磁盘。
此外,所有快照文件都会从 snapshot_local 和 snapshot_s3_plain 移至 snapshot_local2 磁盘。
配置日志缓存
为尽量减少从磁盘读取的数据量,Keeper 会在内存中缓存日志条目。 如果请求较大,日志条目会占用过多内存,因此会限制缓存的日志量。 最新日志缓存的上限由以下参数控制:
latest_logs_cache_size_threshold- 缓存中存储的最新日志总大小
如果默认值过大,可以通过调低此配置来减少内存使用量。
后续提交所需的日志条目由已解码的预读读取器提供,其大小由
log_readahead_commit_window_bytes 控制 (0 会禁用提交预读) 。此设置取代了
已弃用的 commit_logs_cache_size_threshold 和 commit_logs_cache_entry_count_threshold 设置,
后两者仅为保持配置兼容性而保留 (如果未设置前者,它仍会映射到 log_readahead_commit_window_bytes;
后者则不起作用) 。当 log_readahead_enabled 为 true 时,同一预读机制还会为 follower
副本追赶读取提供服务;有关 peer 端的调优参数,请参阅内部协调设置中的
log_readahead_window_bytes、log_readahead_max_peer_readers、log_readahead_eviction_timeout_ms、
log_readahead_pool_threads、log_readahead_serve_wait_timeout_ms 和 log_readahead_chunk_size。
Prometheus
Keeper 可以暴露供 Prometheus 抓取的指标数据。
设置:
endpoint– Prometheus 服务器抓取指标的 HTTP 端点。以 '/' 开头。port–endpoint的端口。metrics– 用于控制是否暴露 system.metrics 表中的指标的标志。events– 用于控制是否暴露 system.events 表中的指标的标志。asynchronous_metrics– 用于控制是否暴露 system.asynchronous_metrics 表中当前指标值的标志。
示例
<clickhouse>
<listen_host>0.0.0.0</listen_host>
<http_port>8123</http_port>
<tcp_port>9000</tcp_port>
<prometheus>
<endpoint>/metrics</endpoint>
<port>9363</port>
<metrics>true</metrics>
<events>true</events>
<asynchronous_metrics>true</asynchronous_metrics>
</prometheus>
</clickhouse>检查 (将 127.0.0.1 替换为你的 ClickHouse 服务器的 IP 地址或主机名) :
curl 127.0.0.1:9363/metrics另请参阅 ClickHouse Cloud Prometheus 集成。
ClickHouse Keeper 用户指南
本指南提供了配置 ClickHouse Keeper 所需的简单最小设置,并附带一个测试分布式操作的示例。该示例在 Linux 上使用 3 个节点进行。
使用 Keeper 设置配置节点
-
在 3 台主机 (
chnode1、chnode2、chnode3) 上安装 3 个 ClickHouse 实例。 (有关 ClickHouse 的安装详情,请参阅快速入门。) -
在每个节点上,添加以下配置项,以允许通过网络接口进行对外通信。
<listen_host>0.0.0.0</listen_host> -
将以下 ClickHouse Keeper 配置添加到三台服务器中的每一台,并为各服务器更新
<server_id>设置;例如,chnode1应设为1,chnode2应设为2,依此类推。<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>warning</raft_logs_level> </coordination_settings> <raft_configuration> <server> <id>1</id> <hostname>chnode1.domain.com</hostname> <port>9234</port> </server> <server> <id>2</id> <hostname>chnode2.domain.com</hostname> <port>9234</port> </server> <server> <id>3</id> <hostname>chnode3.domain.com</hostname> <port>9234</port> </server> </raft_configuration> </keeper_server>以下是上文中使用的基本设置:
参数 描述 示例 tcp_port Keeper 客户端使用的端口 9181,默认对应 ZooKeeper 中的 2181 server_id Raft 配置中每个 ClickHouse Keeper 服务器使用的标识符 1 coordination_settings 协调相关参数部分,例如超时 timeouts: 10000, log level: trace server 参与节点的服务器定义 各服务器定义的列表 raft_configuration Keeper 集群中各服务器的配置 每个服务器的 server 和 settings id Keeper 服务中服务器的数字 ID 1 hostname Keeper 集群中各服务器的主机名、IP 或 FQDN chnode1.domain.comport 用于服务器间 Keeper 连接监听的端口 9234 -
启用 ZooKeeper 组件。它将使用 ClickHouse Keeper 引擎:
<zookeeper> <node> <host>chnode1.domain.com</host> <port>9181</port> </node> <node> <host>chnode2.domain.com</host> <port>9181</port> </node> <node> <host>chnode3.domain.com</host> <port>9181</port> </node> </zookeeper>以下是上文中使用的基本设置:
参数 描述 示例 node 用于 ClickHouse Keeper 连接的节点列表 为每台服务器配置一个 settings 条目 host 每个 ClickHouse Keeper 节点的主机名、IP 或 FQDN chnode1.domain.comport ClickHouse Keeper 客户端端口 9181 -
重启 ClickHouse,并验证每个 Keeper 实例都在运行。在每台服务器上执行以下命令。如果 Keeper 正在运行且状态正常,
ruok命令将返回imok:# echo ruok | nc localhost 9181; echo imok -
system数据库中有一张名为zookeeper的表,包含你的 ClickHouse Keeper 实例的详细信息。我们来查看这张表:SELECT * FROM system.zookeeper WHERE path IN ('/', '/clickhouse')表如下所示:
┌─name───────┬─value─┬─czxid─┬─mzxid─┬───────────────ctime─┬───────────────mtime─┬─version─┬─cversion─┬─aversion─┬─ephemeralOwner─┬─dataLength─┬─numChildren─┬─pzxid─┬─path────────┐ │ clickhouse │ │ 124 │ 124 │ 2022-03-07 00:49:34 │ 2022-03-07 00:49:34 │ 0 │ 2 │ 0 │ 0 │ 0 │ 2 │ 5693 │ / │ │ task_queue │ │ 125 │ 125 │ 2022-03-07 00:49:34 │ 2022-03-07 00:49:34 │ 0 │ 1 │ 0 │ 0 │ 0 │ 1 │ 126 │ /clickhouse │ │ tables │ │ 5693 │ 5693 │ 2022-03-07 00:49:34 │ 2022-03-07 00:49:34 │ 0 │ 3 │ 0 │ 0 │ 0 │ 3 │ 6461 │ /clickhouse │ └────────────┴───────┴───────┴───────┴─────────────────────┴─────────────────────┴─────────┴──────────┴──────────┴────────────────┴────────────┴─────────────┴───────┴─────────────┘
在 ClickHouse 中配置集群
-
让我们在其中 2 个节点上配置一个简单集群:包含 2 个分片,且每个分片只有 1 个副本。第三个节点将用于满足 ClickHouse Keeper 对仲裁的要求。更新
chnode1和chnode2上的配置。以下集群在每个节点上定义了 1 个分片,总共 2 个分片,且不进行复制。在此示例中,一部分数据会位于一个节点上,另一部分会位于另一个节点上:<remote_servers> <cluster_2S_1R> <shard> <replica> <host>chnode1.domain.com</host> <port>9000</port> <user>default</user> <password>ClickHouse123!</password> </replica> </shard> <shard> <replica> <host>chnode2.domain.com</host> <port>9000</port> <user>default</user> <password>ClickHouse123!</password> </replica> </shard> </cluster_2S_1R> </remote_servers>参数 描述 示例 shard 集群定义中的副本列表 每个分片的副本列表 replica 每个副本的设置列表 每个副本的设置项 host 承载副本分片的服务器主机名、IP 或 FQDN chnode1.domain.comport 使用原生 TCP 协议通信时所用的端口 9000 user 用于向集群实例进行身份验证的用户名 default password 为允许连接到集群实例而为该用户设置的密码 ClickHouse123! -
重启 ClickHouse 并验证集群已创建:
SHOW clusters;你应该能看到你的集群:
┌─cluster───────┐ │ cluster_2S_1R │ └───────────────┘
创建并测试分布式表
-
在新集群上,使用
chnode1上的 ClickHouse client 创建一个新数据库。ON CLUSTER子句会自动在两个节点上创建该数据库。CREATE DATABASE db1 ON CLUSTER 'cluster_2S_1R'; -
在
db1数据库中创建一张新表。再次使用ON CLUSTER,会在两个节点上创建该表。CREATE TABLE db1.table1 on cluster 'cluster_2S_1R' ( `id` UInt64, `column1` String ) ENGINE = MergeTree ORDER BY column1 -
在
chnode1节点上添加几行数据:INSERT INTO db1.table1 (id, column1) VALUES (1, 'abc'), (2, 'def') -
在
chnode2节点上添加几行数据:INSERT INTO db1.table1 (id, column1) VALUES (3, 'ghi'), (4, 'jkl') -
请注意,在每个节点上运行
SELECT语句时,只会显示该节点上的数据。例如,在chnode1上:SELECT * FROM db1.table1Query id: 7ef1edbc-df25-462b-a9d4-3fe6f9cb0b6d ┌─id─┬─column1─┐ │ 1 │ abc │ │ 2 │ def │ └────┴─────────┘ 2 rows in set. Elapsed: 0.006 sec.在
chnode2上: -
SELECT * FROM db1.table1Query id: c43763cc-c69c-4bcc-afbe-50e764adfcbf ┌─id─┬─column1─┐ │ 3 │ ghi │ │ 4 │ jkl │ └────┴─────────┘ -
你可以创建一张
Distributed表,用来表示两个分片上的数据。使用Distributed表引擎的表本身不存储任何数据,但支持在多台服务器上进行分布式查询处理。读取会访问所有分片,写入也可以分发到各个分片。在chnode1上运行以下查询:CREATE TABLE db1.dist_table ( id UInt64, column1 String ) ENGINE = Distributed(cluster_2S_1R,db1,table1) -
请注意,查询
dist_table会返回来自这两个分片的全部四行数据:SELECT * FROM db1.dist_tableQuery id: 495bffa0-f849-4a0c-aeea-d7115a54747a ┌─id─┬─column1─┐ │ 1 │ abc │ │ 2 │ def │ └────┴─────────┘ ┌─id─┬─column1─┐ │ 3 │ ghi │ │ 4 │ jkl │ └────┴─────────┘ 4 rows in set. Elapsed: 0.018 sec.
摘要
本指南介绍了如何使用 ClickHouse Keeper 搭建集群。借助 ClickHouse Keeper,您可以配置集群,并定义可在各个分片之间复制的分布式表。
使用唯一路径配置 ClickHouse Keeper
说明
本文介绍如何使用内置的 {uuid} 宏设置,
在 ClickHouse Keeper 或 ZooKeeper 中创建唯一的条目。唯一的
路径 在频繁创建和 drop 表时很有帮助,因为
这样可以避免每次都要等待几分钟,让 Keeper 的垃圾回收
删除 路径 条目;这是因为每次创建 路径 时,都会在该
路径 中使用新的 uuid;路径 永远不会被重复使用。
示例环境
一个三节点集群,配置为在全部三个节点上运行 ClickHouse Keeper, 并在其中两个节点上运行 ClickHouse。这样一来, ClickHouse Keeper 将拥有三个节点 (包括一个仲裁节点) ,并且 ClickHouse 只有一个分片,由两个副本组成。
| node | description |
|---|---|
chnode1.marsnet.local |
数据节点 - 集群 cluster_1S_2R |
chnode2.marsnet.local |
数据节点 - 集群 cluster_1S_2R |
chnode3.marsnet.local |
ClickHouse Keeper 仲裁节点 |
集群配置示例:
<remote_servers>
<cluster_1S_2R>
<shard>
<replica>
<host>chnode1.marsnet.local</host>
<port>9440</port>
<user>default</user>
<password>ClickHouse123!</password>
<secure>1</secure>
</replica>
<replica>
<host>chnode2.marsnet.local</host>
<port>9440</port>
<user>default</user>
<password>ClickHouse123!</password>
<secure>1</secure>
</replica>
</shard>
</cluster_1S_2R>
</remote_servers>将表配置为使用 {uuid} 的步骤
- 在每台服务器上配置宏 服务器 1 的示例:
<macros>
<shard>1</shard>
<replica>replica_1</replica>
</macros>- 创建数据库
CREATE DATABASE db_uuid
ON CLUSTER 'cluster_1S_2R'
ENGINE Atomic;CREATE DATABASE db_uuid ON CLUSTER cluster_1S_2R
ENGINE = Atomic
Query id: 07fb7e65-beb4-4c30-b3ef-bd303e5c42b5
┌─host──────────────────┬─port─┬─status─┬─error─┬─num_hosts_remaining─┬─num_hosts_active─┐
│ chnode2.marsnet.local │ 9440 │ 0 │ │ 1 │ 0 │
│ chnode1.marsnet.local │ 9440 │ 0 │ │ 0 │ 0 │
└───────────────────────┴──────┴────────┴───────┴─────────────────────┴──────────────────┘- 使用宏和
{uuid}在集群中创建表
CREATE TABLE db_uuid.uuid_table1 ON CLUSTER 'cluster_1S_2R'
(
id UInt64,
column1 String
)
ENGINE = ReplicatedMergeTree('/clickhouse/tables/{shard}/db_uuid/{uuid}', '{replica}' )
ORDER BY (id);CREATE TABLE db_uuid.uuid_table1 ON CLUSTER cluster_1S_2R
(
`id` UInt64,
`column1` String
)
ENGINE = ReplicatedMergeTree('/clickhouse/tables/{shard}/db_uuid/{uuid}', '{replica}')
ORDER BY id
Query id: 8f542664-4548-4a02-bd2a-6f2c973d0dc4
┌─host──────────────────┬─port─┬─status─┬─error─┬─num_hosts_remaining─┬─num_hosts_active─┐
│ chnode1.marsnet.local │ 9440 │ 0 │ │ 1 │ 0 │
│ chnode2.marsnet.local │ 9440 │ 0 │ │ 0 │ 0 │
└───────────────────────┴──────┴────────┴───────┴─────────────────────┴──────────────────┘- 创建分布式表
CREATE TABLE db_uuid.dist_uuid_table1 ON CLUSTER 'cluster_1S_2R'
(
id UInt64,
column1 String
)
ENGINE = Distributed('cluster_1S_2R', 'db_uuid', 'uuid_table1' );CREATE TABLE db_uuid.dist_uuid_table1 ON CLUSTER cluster_1S_2R
(
`id` UInt64,
`column1` String
)
ENGINE = Distributed('cluster_1S_2R', 'db_uuid', 'uuid_table1')
Query id: 3bc7f339-ab74-4c7d-a752-1ffe54219c0e
┌─host──────────────────┬─port─┬─status─┬─error─┬─num_hosts_remaining─┬─num_hosts_active─┐
│ chnode2.marsnet.local │ 9440 │ 0 │ │ 1 │ 0 │
│ chnode1.marsnet.local │ 9440 │ 0 │ │ 0 │ 0 │
└───────────────────────┴──────┴────────┴───────┴─────────────────────┴──────────────────┘测试
- 向第一个节点插入数据 (例如
chnode1)
INSERT INTO db_uuid.uuid_table1
( id, column1)
VALUES
( 1, 'abc');INSERT INTO db_uuid.uuid_table1 (id, column1) FORMAT Values
Query id: 0f178db7-50a6-48e2-9a1b-52ed14e6e0f9
Ok.
1 row in set. Elapsed: 0.033 sec.- 向第二个节点插入数据 (例如,
chnode2)
INSERT INTO db_uuid.uuid_table1
( id, column1)
VALUES
( 2, 'def');INSERT INTO db_uuid.uuid_table1 (id, column1) FORMAT Values
Query id: edc6f999-3e7d-40a0-8a29-3137e97e3607
Ok.
1 row in set. Elapsed: 0.529 sec.- 通过分布式表查看记录
SELECT * FROM db_uuid.dist_uuid_table1;SELECT *
FROM db_uuid.dist_uuid_table1
Query id: 6cbab449-9e7f-40fe-b8c2-62d46ba9f5c8
┌─id─┬─column1─┐
│ 1 │ abc │
└────┴─────────┘
┌─id─┬─column1─┐
│ 2 │ def │
└────┴─────────┘
2 rows in set. Elapsed: 0.007 sec.备选方案
默认的复制路径可以预先通过宏定义,也可以使用 {uuid}
- 在每个节点上为表设置默认值
<default_replica_path>/clickhouse/tables/{shard}/db_uuid/{uuid}</default_replica_path>
<default_replica_name>{replica}</default_replica_name>- 创建表时不显式指定参数:
CREATE TABLE db_uuid.uuid_table1 ON CLUSTER 'cluster_1S_2R'
(
id UInt64,
column1 String
)
ENGINE = ReplicatedMergeTree
ORDER BY (id);CREATE TABLE db_uuid.uuid_table1 ON CLUSTER cluster_1S_2R
(
`id` UInt64,
`column1` String
)
ENGINE = ReplicatedMergeTree
ORDER BY id
Query id: ab68cda9-ae41-4d6d-8d3b-20d8255774ee
┌─host──────────────────┬─port─┬─status─┬─error─┬─num_hosts_remaining─┬─num_hosts_active─┐
│ chnode2.marsnet.local │ 9440 │ 0 │ │ 1 │ 0 │
│ chnode1.marsnet.local │ 9440 │ 0 │ │ 0 │ 0 │
└───────────────────────┴──────┴────────┴───────┴─────────────────────┴──────────────────┘
2 rows in set. Elapsed: 1.175 sec.- 验证其是否使用了默认配置中的设置
SHOW CREATE TABLE db_uuid.uuid_table1;SHOW CREATE TABLE db_uuid.uuid_table1
CREATE TABLE db_uuid.uuid_table1
(
`id` UInt64,
`column1` String
)
ENGINE = ReplicatedMergeTree('/clickhouse/tables/{shard}/db_uuid/{uuid}', '{replica}')
ORDER BY id
1 row in set. Elapsed: 0.003 sec.故障排查
获取表信息和 UUID 的示例命令:
SELECT * FROM system.tables
WHERE database = 'db_uuid' AND name = 'uuid_table1';获取 ZooKeeper 中上述表 (含 UUID) 信息的示例命令
SELECT * FROM system.zookeeper
WHERE path = '/clickhouse/tables/1/db_uuid/9e8a3cc2-0dec-4438-81a7-c3e63ce2a1cf/replicas';检查方法如下:
例如,
SELECT name, engine FROM system.databases WHERE name = 'db_uuid';SELECT
name,
engine
FROM system.databases
WHERE name = 'db_uuid'
Query id: b047d459-a1d2-4016-bcf9-3e97e30e49c2
┌─name────┬─engine─┐
│ db_uuid │ Atomic │
└─────────┴────────┘
1 row in set. Elapsed: 0.004 sec.ClickHouse Keeper 动态重新配置
说明
如果启用了 keeper_server.enable_reconfiguration,ClickHouse Keeper 会部分支持 ZooKeeper 的 reconfig
命令,以实现集群的动态重新配置。
虚拟节点 /keeper/config 包含最近一次已提交的集群配置,格式如下:
server.id = server_host:server_port[;server_type][;server_priority]
server.id2 = ...
...- 每个 server 条目以换行符分隔。
server_type的值可以是participant或learner(learner 不参与 leader 选举) 。server_priority是一个非负整数,用于指定在 leader 选举中应优先选择哪些节点。 优先级为 0 表示该 server 永远不会成为 leader。
示例:
:) get /keeper/config
server.1=zoo1:9234;participant;1
server.2=zoo2:9234;participant;1
server.3=zoo3:9234;participant;1你可以使用 reconfig 命令来添加新服务器、移除现有服务器,以及调整现有服务器的优先级。以下是一些示例 (使用 clickhouse-keeper-client) :
# 添加两个新服务器
reconfig add "server.5=localhost:123,server.6=localhost:234;learner"
# 移除另外两个服务器
reconfig remove "3,4"
# 将现有服务器的优先级更改为 8
reconfig add "server.5=localhost:5123;participant;8"以下是 kazoo 的示例:
# 添加两台新服务器,移除另外两台服务器
reconfig(joining="server.5=localhost:123,server.6=localhost:234;learner", leaving="3,4")
# 将现有服务器的优先级更改为 8
reconfig(joining="server.5=localhost:5123;participant;8", leaving=None)joining 中的服务器应采用上文所述的 server 格式。各服务器条目之间应以逗号分隔。
添加新服务器时,可以省略 server_priority (默认值为 1) 和 server_type (默认值
为 participant) 。
如果要更改现有服务器的优先级,请将该服务器及目标优先级一起添加到 joining 中。
服务器主机、端口和类型必须与现有服务器配置相同。
服务器会按照它们在 joining 和 leaving 中出现的顺序添加和移除。
joining 中的所有更新都会先于 leaving 中的更新处理。
Keeper 重新配置的实现有一些注意事项:
-
仅支持增量重新配置。带有非空
new_members的请求会被拒绝。ClickHouse Keeper 的实现依赖 NuRaft API 动态更改成员关系。NuRaft 只能一次 添加一台服务器或移除一台服务器,每次只能处理一个。这意味着对配置的每次更改 (
joining中的每一项、leaving中的每一项) 都必须单独决定。因此不提供批量 重新配置,因为这会对最终用户造成误导。同样也无法更改服务器类型 (participant/learner) ,因为 NuRaft 不支持这一点,而 唯一的办法是先移除再添加服务器,这同样会造成误导。
-
不能使用返回的
znodestat值。 -
from_version字段不会被使用。所有设置了from_version的请求都会被拒绝。 这是因为/keeper/config是一个虚拟节点,也就是说它不会存储在 持久化存储中,而是会针对每个请求根据指定的节点配置动态生成。 作出这一决定是为了避免数据重复,因为 NuRaft 已经存储了该配置。 -
与 ZooKeeper 不同,无法通过提交
sync命令来等待集群重新配置完成。 新配置最终会生效,但不保证具体时间。 -
reconfig命令可能会因各种原因失败。你可以检查集群的状态,查看更新 是否已应用。
将单节点 keeper 转换为集群
有时需要将实验性的 keeper 节点扩展为集群。下面是将其逐步扩展为 3 节点集群的方案:
- 重要:新增节点必须分批添加,且每批数量必须小于当前 法定人数,否则它们会自行选举出一个 leader。本例中按每次一个节点添加。
- 现有 keeper 节点必须开启
keeper_server.enable_reconfiguration配置参数。 - 使用 keeper 集群的完整新配置启动第二个节点。
- 节点启动后,使用
reconfig将其添加到节点 1。 - 然后启动第三个节点,并使用
reconfig将其添加进来。 - 更新
clickhouse-server配置,将新的 keeper 节点添加进去,然后重启以应用更改。 - 更新节点 1 的 raft 配置,并可选择重启它。
为了更好地熟悉这一过程,这里有一个 sandbox 仓库。
不支持的功能
虽然 ClickHouse Keeper 的目标是与 ZooKeeper 完全兼容,但仍有一些功能尚未实现 (不过相关开发仍在持续进行) :
create不支持返回Stat对象create不支持 生存时间 (TTL)addWatch不支持PERSISTENT监听- 不支持
removeWatch和removeAllWatches - 不支持
setWatches - 不支持创建
CONTAINER类型的 znode - 不支持
SASL 身份验证