Skip to content
ClickHouse Docs
ClickHouse DocsClickHouse Docs

ClickHouse Keeper

Не поддерживается в ClickHouse Cloud

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" для получения подробной информации о состоянии с точки зрения кворума и о клиентских подключениях.
imok
  • mntr: Выводит список переменных, которые можно использовать для мониторинга состояния кластера.
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     0

zk_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: 4
  • stat: Выводит краткие сведения о сервере и подключенных клиентах.
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: 4
  • srst: Сбрасывает статистику сервера. Команда повлияет на результаты 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=20
  • cons: Показать полные сведения обо всех соединениях/сеансах клиентов, подключенных к этому серверу. Включает информацию о количестве полученных/отправленных пакетов, идентификаторе сеанса, задержках операций, последней выполненной операции и т. д…
 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: 3875
  • isro: Проверяет, работает ли сервер в режиме только для чтения. Сервер вернёт ro, если включён режим только для чтения, и rw, если нет.
rw
  • wchs: Выводит краткую информацию о наблюдениях на сервере.
1 connections watching 1 paths
Total watches:1
  • wchc: Выводит подробную информацию о наблюдениях на сервере по сеансам. В результате отображается список сеансов (соединений) со связанными наблюдениями (путями). Обратите внимание: в зависимости от количества наблюдений эта операция может быть ресурсоёмкой (влиять на производительность сервера), поэтому используйте её с осторожностью.
0x0000000000000001
    /clickhouse/task_queue/ddl
  • wchp: Выводит подробную информацию о наблюдениях на сервере по путям. В результате выводится список путей (znode) со связанными сеансами. Обратите внимание: в зависимости от количества наблюдений эта операция может быть ресурсоёмкой (то есть влиять на производительность сервера), поэтому используйте её с осторожностью.
/clickhouse/task_queue/ddl
    0x0000000000000001
  • dump: Выводит активные сеансы и эфемерные узлы. Команда работает только на лидере.
Sessions dump (2):
0x0000000000000001
0x0000000000000002
Sessions with Ephemerals (1):
0x0000000000000001
 /clickhouse/task_queue/ddl
  • csnp: Планирует задачу создания снимка. В случае успеха возвращает последний зафиксированный индекс журнала для запланированного снимка, а в случае ошибки — Failed to schedule snapshot creation task.. Команда lgif поможет определить, завершено ли создание снимка.
100
  • lgif: информация о журнале 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   50
  • rqld: Запрос на назначение новым лидером. Возвращает 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    0
  • ydld: Запрашивает отказ от лидерства и переход в состояние фолловера. Если сервер, получивший запрос, является лидером, он сначала приостановит операции записи, дождётся, пока преемник (текущий лидер не может быть преемником) завершит дозагрузку последнего журнала, а затем сложит полномочия. Преемник будет выбран автоматически. Возвращает 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;

Зафиксируйте метрики для сравнения перед миграцией, чтобы затем проверить согласованность.

Этапы миграции

  1. Остановите ингестию данных на всех узлах ClickHouse.

  2. Остановите все фоновые задачи на всех узлах ClickHouse (см. выше).

  3. Остановите все узлы ZooKeeper.

  4. Необязательно, но рекомендуется: найдите узел-лидер ZooKeeper, запустите его и затем снова остановите. Это принудительно заставит ZooKeeper записать на диск согласованный снимок перед преобразованием.

  5. Запустите 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
  1. Скопируйте снимок на все узлы ClickHouse Keeper. Снимок должен находиться на каждом узле до запуска любого из них — если узел запустится без снимка, он может избрать себя лидером с пустым состоянием.

  2. Обновите конфигурацию ClickHouse, чтобы она указывала на новый кластер ClickHouse Keeper.

  3. Запустите ClickHouse Keeper на всех узлах, затем перезапустите ClickHouse.

  4. Сравните метрики с базовыми значениями до миграции, чтобы убедиться в согласованности.

  5. Возобновите фоновые задачи и перезапустите ингестию данных.

Объединение нескольких кластеров ZooKeeper

Если вы используете несколько кластеров ZooKeeper — например, по одному на каждую группу сегментов, — их можно объединить в один кластер ClickHouse Keeper. Официальный инструмент clickhouse-keeper-converter поддерживает только конвертацию в формате «один к одному» (один кластер ZooKeeper в один снимок Keeper), поэтому для консолидации потребуется изменить исходный код конвертера, чтобы объединить несколько снимков:

  1. Запустите clickhouse-keeper-converter отдельно для каждого кластера ZooKeeper, записывая каждый результат в отдельный каталог.
  2. Последовательно десериализуйте файлы снимков. При слиянии пересчитайте значения numChildren, чтобы избежать конфликтов идентификаторов узлов между пространствами имен из разных исходных кластеров.
  3. Запишите объединенный результат в целевой каталог снимков 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 узел. Использовать его следует только в крайнем случае: если вы не можете снова запустить свои узлы или запустить новый экземпляр на той же конечной точке.

Перед продолжением обратите внимание на следующее:

  • Убедитесь, что вышедшие из строя узлы больше не смогут подключиться к кластеру.
  • Не запускайте новые узлы, пока это не будет указано в шагах.

Убедившись, что эти условия выполнены, сделайте следующее:

  1. Выберите один узел Keeper, который станет новым лидером. Учтите, что данные именно этого узла будут использованы для всего кластера, поэтому мы рекомендуем выбрать узел с наиболее актуальным состоянием.
  2. Прежде чем делать что-либо еще, создайте резервную копию каталогов log_storage_path и snapshot_storage_path на выбранном узле.
  3. Выполните реконфигурацию кластера на всех узлах, которые вы хотите использовать.
  4. Отправьте на выбранный узел четырехбуквенную команду rcvr, чтобы перевести узел в режим восстановления, ИЛИ остановите экземпляр Keeper на выбранном узле и снова запустите его с аргументом --force-recovery.
  5. По одному запускайте экземпляры Keeper на новых узлах, убеждаясь, что mntr возвращает follower для zk_server_state перед запуском следующего узла.
  6. В режиме восстановления узел-лидер будет возвращать сообщение об ошибке для команды mntr, пока не достигнет кворума с новыми узлами, и будет отклонять любые запросы от клиента и последователей.
  7. После достижения кворума узел-лидер вернется в обычный режим работы и начнет принимать все запросы с проверкой через 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

  1. Установите 3 экземпляра ClickHouse на 3 хостах (chnode1, chnode2, chnode3). (Подробности по установке ClickHouse см. в разделе Быстрый старт.)

  2. На каждом узле добавьте следующую запись, чтобы разрешить внешние подключения через сетевой интерфейс.

    <listen_host>0.0.0.0</listen_host>
  3. Добавьте следующую конфигурацию ClickHouse Keeper на все три сервера, изменив параметр <server_id> для каждого из них; для chnode1 это будет 1, для chnode22 и т. д.

    <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.com
    port порт для межсерверных соединений Keeper 9234
  4. Включите компонент 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.com
    port клиентский порт ClickHouse Keeper 9181
  5. Перезапустите ClickHouse и убедитесь, что каждый экземпляр Keeper запущен. Выполните следующую команду на каждом сервере. Команда ruok возвращает imok, если Keeper запущен и работает нормально:

    # echo ruok | nc localhost 9181; echo
    imok
  6. В базе данных 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

  1. Давайте настроим простой кластер с 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.com
    port порт, используемый для взаимодействия через собственный TCP-протокол 9000
    user имя пользователя, которое будет использоваться для аутентификации при подключении к экземплярам кластера default
    password пароль пользователя, которому разрешены подключения к экземплярам кластера ClickHouse123!
  2. Перезапустите ClickHouse и убедитесь, что кластер создан:

    SHOW clusters;

    Вы должны увидеть свой кластер:

    ┌─cluster───────┐
    │ cluster_2S_1R │
    └───────────────┘

Создайте и протестируйте distributed таблицу

  1. Создайте новую базу данных в новом кластере с помощью клиента ClickHouse на chnode1. Предложение ON CLUSTER автоматически создаёт базу данных на обоих узлах.

    CREATE DATABASE db1 ON CLUSTER 'cluster_2S_1R';
  2. Создайте новую таблицу в базе данных db1. И снова ON CLUSTER создаёт таблицу на обоих узлах.

    CREATE TABLE db1.table1 on cluster 'cluster_2S_1R'
    (
        `id` UInt64,
        `column1` String
    )
    ENGINE = MergeTree
    ORDER BY column1
  3. На узле chnode1 добавьте пару строк:

    INSERT INTO db1.table1
        (id, column1)
    VALUES
        (1, 'abc'),
        (2, 'def')
  4. Добавьте пару строк на узле chnode2:

    INSERT INTO db1.table1
        (id, column1)
    VALUES
        (3, 'ghi'),
        (4, 'jkl')
  5. Обратите внимание, что выполнение оператора SELECT на каждом узле показывает только данные, находящиеся на этом узле. Например, на chnode1:

    SELECT *
    FROM db1.table1
    Query id: 7ef1edbc-df25-462b-a9d4-3fe6f9cb0b6d
    
    ┌─id─┬─column1─┐
    │  1 │ abc     │
    │  2 │ def     │
    └────┴─────────┘
    
    2 rows in set. Elapsed: 0.006 sec.

    На chnode2:

  6. SELECT *
    FROM db1.table1
    Query id: c43763cc-c69c-4bcc-afbe-50e764adfcbf
    
    ┌─id─┬─column1─┐
    │  3 │ ghi     │
    │  4 │ jkl     │
    └────┴─────────┘
  7. Вы можете создать таблицу Distributed, чтобы представить данные на двух сегментах. Таблицы с движком таблицы Distributed не хранят собственные данные, но позволяют выполнять распределённую обработку запросов на нескольких серверах. Чтение выполняется со всех сегментов, а операции записи можно распределять между сегментами. Выполните следующий запрос на chnode1:

    CREATE TABLE db1.dist_table (
        id UInt64,
        column1 String
    )
    ENGINE = Distributed(cluster_2S_1R,db1,table1)
  8. Обратите внимание, что запрос к dist_table возвращает все четыре строки данных из двух сегментов:

    SELECT *
    FROM db1.dist_table
    Query 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 с уникальными путями

Не поддерживается в ClickHouse Cloud

Описание

В этой статье описано, как использовать встроенную настройку макроса {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. Настройте макросы на каждом сервере пример для сервера 1:
    <macros>
        <shard>1</shard>
        <replica>replica_1</replica>
    </macros>
  1. Создайте базу данных
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 │
└───────────────────────┴──────┴────────┴───────┴─────────────────────┴──────────────────┘
  1. Создайте таблицу в кластере с помощью макросов и {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 │
└───────────────────────┴──────┴────────┴───────┴─────────────────────┴──────────────────┘
  1. Создайте 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 │
└───────────────────────┴──────┴────────┴───────┴─────────────────────┴──────────────────┘

Тестирование

  1. Вставьте данные в первый узел (например, 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.
  1. Вставьте данные на второй узел (например, 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.
  1. Просмотрите записи с использованием 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}

  1. Задайте значение по умолчанию для таблиц на каждом узле
<default_replica_path>/clickhouse/tables/{shard}/db_uuid/{uuid}</default_replica_path>
<default_replica_name>{replica}</default_replica_name>
  1. Создайте таблицу без явного указания параметров:
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.
  1. Убедитесь, что используются настройки из конфигурации по умолчанию
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 Cloud

Описание

ClickHouse Keeper частично поддерживает команду ZooKeeper reconfig для динамической реконфигурации кластера, если включён параметр keeper_server.enable_reconfiguration.

Виртуальный узел /keeper/config содержит последнюю подтверждённую конфигурацию кластера в следующем формате:

server.id = server_host:server_port[;server_type][;server_priority]
server.id2 = ...
...

Пример:

:) 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, некоторые возможности пока не реализованы (разработка продолжается):

Navigation