ClickHouse Keeper обеспечивает систему координации для репликации данных и выполнения запросов distributed DDL. ClickHouse Keeper совместим с ZooKeeper.
Подробности реализации
ZooKeeper — одна из первых широко известных систем координации с открытым исходным кодом. Она реализована на Java и использует простую, но мощную модель данных. Алгоритм координации ZooKeeper, ZooKeeper Atomic Broadcast (ZAB), не дает гарантий линеаризуемости для операций чтения, поскольку каждый узел ZooKeeper обрабатывает чтение локально. В отличие от ZooKeeper, ClickHouse Keeper написан на C++ и использует реализацию алгоритма RAFT. Этот алгоритм обеспечивает линеаризуемость операций чтения и записи и имеет несколько реализаций с открытым исходным кодом на разных языках.
По умолчанию ClickHouse Keeper предоставляет те же гарантии, что и ZooKeeper: линеаризуемые записи и нелинеаризуемые чтения. Он использует совместимый клиент-серверный протокол, поэтому для работы с ClickHouse Keeper можно использовать любой стандартный клиент ZooKeeper. Снимки и журналы имеют формат, несовместимый с 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. В обоих случаях используется практически один и тот же файл конфигурации .xml.
Настройки конфигурации Keeper
Основной тег конфигурации ClickHouse Keeper — <keeper_server>. Он имеет следующие параметры:
| Parameter | Description | Default |
|---|---|---|
tcp_port |
Порт для подключения клиента. | 2181 |
tcp_port_secure |
Защищенный порт для SSL-соединения между клиентом и сервером Keeper. | - |
server_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, этот интервал используется для отслеживания объема физической памяти. При изменении объема памяти мягкий лимит памяти Keeper пересчитывается с использованием max_memory_usage_soft_limit_ratio. |
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 отправляет heartbeat-сообщения узлам follower (мс) | 500 |
election_timeout_lower_bound_ms |
Если узел follower не получает heartbeat от лидера в течение этого интервала, он может инициировать выборы лидера. Значение должно быть меньше или равно election_timeout_upper_bound_ms. В идеале эти значения не должны быть равны. |
1000 |
election_timeout_upper_bound_ms |
Если узел follower не получает heartbeat от лидера в течение этого интервала, он обязан инициировать выборы лидера. | 2000 |
rotate_log_storage_interval |
Сколько записей журнала хранить в одном файле. | 100000 |
reserved_log_items |
Сколько записей журнала координации хранить до компактизации. | 100000 |
snapshot_distance |
Как часто ClickHouse Keeper создает новые снимки (по количеству записей в журналах). | 100000 |
snapshots_to_keep |
Сколько снимков хранить. | 3 |
stale_log_gap |
Порог, при котором лидер считает узел follower устаревшим и отправляет ему снимок вместо журналов. | 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 |
Разрешить пересылку запросов на запись от узлов follower лидеру. | true |
shutdown_timeout |
Сколько ждать завершения внутренних соединений и остановки (мс). | 5000 |
startup_timeout |
Если сервер не подключится к другим участникам кворума за указанный тайм-аут, он завершит работу (мс). | 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 при дозагрузке журнала изменений (репликация follower). | false |
log_readahead_window_bytes |
Максимальное число байтов декодированных записей, буферизуемых для каждого считывателя peer. Должно быть не меньше размера типичного батча append-entries. | 64MiB |
log_readahead_max_peer_readers |
Максимальное число одновременных считывателей с упреждающей загрузкой для каждого peer. | 8 |
log_readahead_eviction_timeout_ms |
Тайм-аут простоя, по истечении которого неактивный считыватель с упреждающей загрузкой для каждого peer или считыватель commit вытесняется (при следующем промахе он создается заново). | 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 |
Максимальный общий размер декодированных записей журнала, буферизуемых для потока commit. 0 отключает предварительное чтение для commit (commit читает записи с диска по одной). |
500MiB |
log_startup_read_max_streams |
Максимальное число файлов журнала изменений, читаемых одновременно во время запуска Keeper. 0 = автоматически использовать число ядер CPU. 1 = использовать последовательное чтение при запуске (до распараллеливания). Фактический параллелизм ограничен числом файлов журнала изменений, которые необходимо прочитать; рассмотрите снижение значения для хранилищ, ограниченных операциями seek (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 в качестве backend-хранилища | 0 |
Конфигурация кворума находится в разделе <keeper_server>.<raft_configuration> и содержит описание серверов.
Единственный параметр для всего кворума — secure, который включает шифрование соединения для взаимодействия между участниками кворума. Установите параметр в true, если для внутреннего взаимодействия между узлами требуется SSL-соединение, в противном случае не указывайте его.
Основные параметры для каждого <server>:
id— Идентификатор сервера в кворуме.hostname— Имя хоста, на котором расположен этот сервер.port— Порт, на котором этот сервер прослушивает подключения.can_become_leader— Установитеfalse, чтобы настроить сервер какlearner. Если параметр не указан, используется значениеtrue.
Примеры конфигурации кворума из трёх узлов можно найти в integration tests с префиксом test_keeper_. Пример конфигурации для сервера №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.
Вы можете отправлять команды в ClickHouse Keeper через telnet или nc, используя клиентский порт.
echo mntr | nc localhost 9181Ниже приведены подробные команды 4lw:
ruok: Проверяет, работает ли сервер и нет ли у него ошибок. Если сервер работает, он ответитimok. В противном случае он не ответит вовсе. Ответimokне обязательно означает, что сервер присоединился к кворуму, а лишь то, что процесс сервера активен и привязан к указанному клиентскому порту. Используйте "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 — накопительные метрики, доступные только лидеру. Это наблюдения на уровне отдельного сервера, а не показатели доступности всего кластера: отсчёт начинается, когда сервер перестаёт наблюдать активного лидера. Во время сетевой партиции сюда может входить время, в течение которого прежний лидер остаётся активным в другой партиции. Метрики выборов фиксируют только успешные выборы, следующие за таким локально наблюдаемым окном отсутствия лидера; передача лидерства, при которой не фиксируется состояние отсутствия лидера при опросе, намеренно не учитывается.
Keeper считывает локальное состояние лидера NuRaft раз в heart_beat_interval_ms, но не чаще одного раза в 100 миллисекунд. Граница каждого окна может отстоять от соответствующего локального перехода состояния не более чем на фактический интервал, а окна короче этого интервала могут быть пропущены. Keeper фиксирует локальное завершение выборов при NuRaft BecomeLeader. srst сбрасывает все четыре значения.
Длительность последнего локально наблюдаемого окна каждого типа также экспортируется в system.asynchronous_metrics как KeeperLastLeaderElectionTime и KeeperLastLeaderUnavailableTime в миллисекундах. Для KeeperLastLeaderElectionTime используется то же определение окна отсутствия лидера, что и для zk_sum_election_time. Обе метрики также доступны только лидеру: узел, не являющийся активным лидером или ещё не завершивший окно, возвращает 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: Показать полные сведения обо всех соединениях/сеансах клиентов, подключенных к этому серверу. Включает информацию о количестве полученных/отправленных пакетов, идентификаторе сеанса, задержках операций, последней выполненной операции и т. д…
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: Проверяет, работает ли сервер в режиме только для чтения. Сервер вернётro, если включён режим только для чтения, иrw, если нет.
rwwchs: Выводит краткую информацию о наблюдениях на сервере.
1 connections watching 1 paths
Total watches:1wchc: Выводит подробную информацию о наблюдениях на сервере по сеансам. В результате отображается список сеансов (соединений) со связанными наблюдениями (путями). Обратите внимание: в зависимости от количества наблюдений эта операция может быть ресурсоёмкой (влиять на производительность сервера), поэтому используйте её с осторожностью.
0x0000000000000001
/clickhouse/task_queue/ddlwchp: Выводит подробную информацию о наблюдениях на сервере по путям. В результате выводится список путей (znode) со связанными сеансами. Обратите внимание: в зависимости от количества наблюдений эта операция может быть ресурсоёмкой (то есть влиять на производительность сервера), поэтому используйте её с осторожностью.
/clickhouse/task_queue/ddl
0x0000000000000001dump: Выводит активные сеансы и эфемерные узлы. Команда работает только на лидере.
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: зафиксированный индекс журнала лидера с моей точки зрения;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: Запрос на назначение новым лидером. ВозвращаетSent leadership request to leader., если запрос отправлен, илиFailed to send leadership request to leader., если запрос не отправлен. Если узел уже является лидером, результат будет таким же, как при успешной отправке запроса.
Sent leadership request to leader.ftfl: Выводит все флаги возможностей и указывает, включены ли они для экземпляра Keeper.
filtered_list 1
multi_read 1
check_not_exists 0ydld: Запрашивает отказ от лидерства и переход в состояние фолловера. Если сервер, получивший запрос, является лидером, он сначала приостановит операции записи, дождётся, пока преемник (текущий лидер не может быть преемником) завершит дозагрузку последнего журнала, а затем сложит полномочия. Преемник будет выбран автоматически. ВозвращаетSent yield leadership request to leader., если запрос отправлен, илиFailed to send yield leadership request to leader., если запрос не отправлен. Если узел уже является фолловером, результат будет таким же, как при отправке запроса.
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 |
Поддержка запроса на чтение нескольких значений | 1 |
filtered_list |
Поддержка запроса списка, который фильтрует результаты по типу узла (эфемерный или постоянный) |
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, запустите его и затем снова остановите. Это принудительно заставит ZooKeeper записать на диск согласованный снимок перед преобразованием.
-
Запустите
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. Снимок должен находиться на каждом узле до запуска любого из них — если узел запустится без снимка, он может избрать себя лидером с пустым состоянием.
-
Обновите конфигурацию ClickHouse, чтобы она указывала на новый кластер ClickHouse Keeper.
-
Запустите ClickHouse Keeper на всех узлах, затем перезапустите ClickHouse.
-
Сравните метрики с базовыми значениями до миграции, чтобы убедиться в согласованности.
-
Возобновите фоновые задачи и перезапустите ингестию данных.
Объединение нескольких кластеров ZooKeeper
Если вы используете несколько кластеров ZooKeeper — например, по одному на каждую группу сегментов, — их можно объединить в один кластер ClickHouse Keeper. Официальный инструмент clickhouse-keeper-converter поддерживает только конвертацию в формате «один к одному» (один кластер ZooKeeper в один снимок Keeper), поэтому для консолидации потребуется изменить исходный код конвертера, чтобы объединить несколько снимков:
- Запустите
clickhouse-keeper-converterотдельно для каждого кластера ZooKeeper, записывая каждый результат в отдельный каталог. - Последовательно десериализуйте файлы снимков. При слиянии пересчитайте значения
numChildren, чтобы избежать конфликтов идентификаторов узлов между пространствами имен из разных исходных кластеров. - Запишите объединенный результат в целевой каталог снимков ClickHouse Keeper.
Работа с шифрованием и ACL
ClickHouse Keeper поддерживает те же схемы ACL, что и ZooKeeper (world, auth, digest). Способ обработки ACL при конвертации зависит от вашей конфигурации ZooKeeper:
- Полностью зашифровано или полностью незашифровано: Конвертируйте напрямую. Конвертер сохранит существующую информацию об ACL.
- Частично зашифровано: Перед конвертацией выдайте права учётной записи суперадминистратора и очистите ACL с помощью
setAcl -Rдля затронутых путей. Выполните конвертацию, затем при необходимости снова включите шифрование в ClickHouse Keeper.
Проверка миграции
После запуска ClickHouse Keeper и перезапуска ClickHouse сравните ключевые метрики с базовыми значениями, зафиксированными до миграции, чтобы убедиться, что миграция прошла успешно.
При консолидации нескольких кластеров ZooKeeper различайте:
- Общие пути: пути, присутствующие в нескольких исходных кластерах и содержащие одинаковые данные, — для них в объединённом результате следует выполнить дедупликацию.
- Различающиеся пути: пути, существующие только в определённых кластерах (например, в
/clickhouse/tablesдля каждой группы сегментов), — их необходимо сохранить из соответствующего источника.
Не обходите большие деревья ZooKeeper напрямую для сравнения. Вместо этого в процессе преобразования выводите все преобразованные пути в файл.
Тонкая настройка после миграции
После миграции при необходимости настройте следующие параметры для более крупных кластеров или более высокой пропускной способности:
| Setting | Default | Recommended | Notes |
|---|---|---|---|
max_requests_batch_size |
100 |
10000 |
Увеличьте для кластеров с большим количеством частей или сегментов |
force_sync |
true |
false |
Асинхронная запись в журнал повышает пропускную способность |
compress_logs |
false |
true |
Сжимает файлы журнала Raft, уменьшая дисковый I/O |
compress_snapshots_with_zstd_format |
true |
— | Уже включен по умолчанию; сжимает снимки с помощью zstd |
Эти параметры задаются в разделе coordination_settings вашей конфигурации Keeper.
Восстановление после потери кворума
Поскольку ClickHouse Keeper использует Raft, он может выдержать определенное количество отказов узлов в зависимости от размера кластера. Например, в кластере из 3 узлов он продолжит корректно работать, если выйдет из строя только 1 узел.
Конфигурацию кластера можно менять динамически, но есть некоторые ограничения. Реконфигурация также опирается на Raft, поэтому для добавления или удаления узла из кластера необходим кворум. Если одновременно потерять слишком много узлов кластера и не будет возможности снова их запустить, Raft перестанет работать и не позволит реконфигурировать кластер обычным способом.
Тем не менее, в ClickHouse Keeper есть режим восстановления, который позволяет принудительно реконфигурировать кластер, имея только 1 узел. Использовать его следует только в крайнем случае: если вы не можете снова запустить свои узлы или запустить новый экземпляр на той же конечной точке.
Перед продолжением обратите внимание на следующее:
- Убедитесь, что вышедшие из строя узлы больше не смогут подключиться к кластеру.
- Не запускайте новые узлы, пока это не будет указано в шагах.
Убедившись, что эти условия выполнены, сделайте следующее:
- Выберите один узел Keeper, который станет новым лидером. Учтите, что данные именно этого узла будут использованы для всего кластера, поэтому мы рекомендуем выбрать узел с наиболее актуальным состоянием.
- Прежде чем делать что-либо еще, создайте резервную копию каталогов
log_storage_pathиsnapshot_storage_pathна выбранном узле. - Выполните реконфигурацию кластера на всех узлах, которые вы хотите использовать.
- Отправьте на выбранный узел четырехбуквенную команду
rcvr, чтобы перевести узел в режим восстановления, ИЛИ остановите экземпляр Keeper на выбранном узле и снова запустите его с аргументом--force-recovery. - По одному запускайте экземпляры Keeper на новых узлах, убеждаясь, что
mntrвозвращаетfollowerдляzk_server_stateперед запуском следующего узла. - В режиме восстановления узел-лидер будет возвращать сообщение об ошибке для команды
mntr, пока не достигнет кворума с новыми узлами, и будет отклонять любые запросы от клиента и последователей. - После достижения кворума узел-лидер вернется в обычный режим работы и начнет принимать все запросы с проверкой через Raft — проверьте это с помощью
mntr, который должен возвращатьleaderдляzk_server_state.
Использование дисков с 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, log_storage_disk должен быть диском local, если latest_log_storage_disk не используется.
Если используется 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 при запуске попытается автоматически переместить файлы на нужные диски. Действует та же гарантия, что и раньше: пока файл не будет полностью перемещён на новый диск, он не удаляется со старого, поэтому можно безопасно выполнять несколько перезапусков.
Если требуется переместить файлы на полностью новый диск (или перейти с конфигурации с 2 дисками на конфигурацию с одним диском), можно использовать несколько определений keeper_server.old_snapshot_storage_disk и keeper_server.old_log_storage_disk.
Следующая конфигурация показывает, как перейти с предыдущей конфигурации с 2 дисками на полностью новую конфигурацию с одним диском:
<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; параметры
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
для настройки на стороне peer см. во внутренних настройках координации.
Prometheus
Keeper может предоставлять метрики для сбора Prometheus.
Настройки:
endpoint– HTTP-конечная точка для сбора метрик сервером Prometheus. Должна начинаться с '/'.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 на IP-адрес или имя хоста вашего сервера ClickHouse):
curl 127.0.0.1:9363/metricsСм. также интеграцию Prometheus для ClickHouse Cloud.
Руководство пользователя ClickHouse Keeper
В этом руководстве приведены простые и минимально необходимые настройки для конфигурации ClickHouse Keeper, а также пример проверки распределённых операций. В примере используются 3 узла Linux.
Настройте узлы, указав параметры Keeper
-
Установите 3 экземпляра ClickHouse на 3 хостах (
chnode1,chnode2,chnode3). (Подробности по установке 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, эквивалент порта 2181 по умолчанию, как в ZooKeeper server_id идентификатор каждого сервера ClickHouse Keeper, используемый в конфигурации Raft 1 coordination_settings раздел с параметрами, такими как тайм-ауты тайм-ауты: 10000, уровень логирования: trace server описание участвующего сервера список описаний всех серверов raft_configuration настройки каждого сервера в кластере Keeper сервер и настройки для каждого id числовой идентификатор сервера для сервисов Keeper 1 hostname имя хоста, IP-адрес или FQDN каждого сервера в кластере Keeper 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 запись в настройках для каждого сервера host имя хоста, IP-адрес или FQDN каждого узла ClickHouse Keeper chnode1.domain.comport клиентский порт ClickHouse Keeper 9181 -
Перезапустите ClickHouse и убедитесь, что каждый экземпляр Keeper запущен. Выполните следующую команду на каждом сервере. Команда
ruokвозвращаетimok, если Keeper запущен и работает нормально:# 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 сегментами и только одной репликой на двух узлах. Третий узел будет использоваться для обеспечения кворума, необходимого 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 │ └───────────────┘
Создайте и протестируйте distributed таблицу
-
Создайте новую базу данных в новом кластере с помощью клиента ClickHouse на
chnode1. Предложение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 можно настраивать кластеры и определять distributed таблицы, которые могут реплицироваться между сегментами.
Настройка ClickHouse Keeper с уникальными путями
Описание
В этой статье описано, как использовать встроенную настройку макроса {uuid}
для создания уникальных записей в ClickHouse Keeper или ZooKeeper. Уникальные
пути полезны при частом создании и удалении таблиц, поскольку
позволяют не ждать несколько минут, пока сборщик мусора Keeper
удалит записи путей: при каждом создании пути в нём используется новый uuid,
поэтому пути никогда не используются повторно.
Пример окружения
Кластер из трёх узлов будет настроен так, что ClickHouse Keeper будет работать на всех трёх узлах, а ClickHouse — на двух из них. Это обеспечивает для ClickHouse Keeper три узла (включая узел для разрешения спорных ситуаций) и один сегмент ClickHouse, состоящий из двух реплик.
| node | description |
|---|---|
chnode1.marsnet.local |
узел данных — cluster cluster_1S_2R |
chnode2.marsnet.local |
узел данных — cluster 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 │
└───────────────────────┴──────┴────────┴───────┴─────────────────────┴──────────────────┘- Создайте distributed таблицу
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.- Просмотрите записи с использованием distributed таблицы
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
Описание
ClickHouse Keeper частично поддерживает команду ZooKeeper reconfig
для динамической реконфигурации кластера, если включён параметр keeper_server.enable_reconfiguration.
Виртуальный узел /keeper/config содержит последнюю подтверждённую конфигурацию кластера в следующем формате:
server.id = server_host:server_port[;server_type][;server_priority]
server.id2 = ...
...- Каждая запись о server отделяется новой строкой.
server_type— либоparticipant, либоlearner(learner не участвует в выборах лидера).server_priority— неотрицательное целое число, указывающее, каким узлам следует отдавать приоритет при выборе лидера. Приоритет 0 означает, что server никогда не станет лидером.
Пример:
:) 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 использует API NuRaft для динамического изменения состава участников. NuRaft позволяет добавлять или удалять только по одному серверу за раз. Это означает, что каждое изменение конфигурации (каждая часть
joining, каждая частьleaving) должно рассматриваться отдельно. Поэтому массовая реконфигурация недоступна, поскольку это могло бы ввести конечных пользователей в заблуждение.Изменить тип сервера (participant/learner) тоже невозможно, поскольку это не поддерживается NuRaft, а единственный вариант — удалить сервер и добавить его заново, что, опять же, вводило бы в заблуждение.
-
Нельзя использовать возвращаемое значение
znodestat. -
Поле
from_versionне используется. Все запросы с заданнымfrom_versionотклоняются. Это связано с тем, что/keeper/config— виртуальный узел, то есть он не хранится в постоянном хранилище, а генерируется на лету на основе указанной конфигурации узла для каждого запроса. Такое решение было принято, чтобы не дублировать данные, поскольку NuRaft уже хранит эту конфигурацию. -
В отличие от ZooKeeper, нельзя дождаться завершения реконфигурации кластера, отправив команду
sync. Новая конфигурация в конечном итоге будет применена, но без каких-либо гарантий по времени. -
Команда
reconfigможет завершиться ошибкой по разным причинам. Вы можете проверить состояние кластера и посмотреть, было ли применено обновление.
Преобразование одноузлового Keeper в кластер
Иногда возникает необходимость расширить экспериментальный узел Keeper до кластера. Ниже приведена пошаговая схема для кластера из 3 узлов:
- ВАЖНО: новые узлы нужно добавлять батчами, размер которых меньше текущего кворума, иначе они выберут лидера между собой. В этом примере — по одному.
- На существующем узле Keeper должен быть включён параметр конфигурации
keeper_server.enable_reconfiguration. - Запустите второй узел с полной новой конфигурацией кластера Keeper.
- После запуска добавьте его на узел 1 с помощью
reconfig. - Затем запустите третий узел и добавьте его с помощью
reconfig. - Обновите конфигурацию
clickhouse-server, добавив в неё новый узел Keeper, и перезапустите его, чтобы применить изменения. - Обновите конфигурацию Raft на узле 1 и, при необходимости, перезапустите его.
Чтобы лучше понять процесс, вот репозиторий-песочница.
Неподдерживаемые возможности
Хотя ClickHouse Keeper стремится к полной совместимости с ZooKeeper, некоторые возможности пока не реализованы (разработка продолжается):
createне поддерживает возврат объектаStatcreateне поддерживает TTLaddWatchне работает с наблюдениямиPERSISTENTremoveWatchиremoveAllWatchesне поддерживаютсяsetWatchesне поддерживается- Создание znode типа
CONTAINERне поддерживается SASL authenticationне поддерживается