В режиме KRaft метаданные Apache Kafka хранятся в собственном журнале, а за их согласование отвечает quorum контроллеров. Для рабочего кластера роли рекомендуется разделять: узлы controller занимаются метаданными, а узлы broker принимают клиентские подключения и хранят разделы топиков.
Ниже приведён пример нового кластера Apache Kafka KRaft с динамическим quorum. Мы создадим три controller и три broker. Миграция существующей установки с ZooKeeper, настройка TLS, SASL и ACL в эту инструкцию не входят.
Как будет устроен кластер
Для примера используем шесть VDS в одной приватной сети:
controller-1.internal—10.20.0.11,node.id=1001;controller-2.internal—10.20.0.12,node.id=1002;controller-3.internal—10.20.0.13,node.id=1003;broker-1.internal—10.20.0.21,node.id=2001;broker-2.internal—10.20.0.22,node.id=2002;broker-3.internal—10.20.0.23,node.id=2003.
Контроллеры будут принимать KRaft-соединения на TCP-порту 9093, брокеры — межброкерные и клиентские соединения на TCP-порту 9092. Имена, адреса и номера портов можно изменить, но они должны быть согласованы во всех конфигурациях и доступны по сети между нужными узлами.
Для metadata quorum требуется большинство работающих контроллеров. Quorum из трёх узлов продолжает работать при отказе одного из них. Пять контроллеров выдерживают отказ двух. Чётное количество не даёт соответствующего выигрыша: например, quorum из четырёх узлов также требует трёх доступных контроллеров.
Размещение трёх процессов controller на одном физическом сервере или одной VDS не обеспечивает отказоустойчивость. Контроллеры должны находиться на независимых узлах, а по возможности — в разных группах отказа инфраструктуры.

Для такой схемы нужны административный доступ к серверам, изолированная сеть и достаточные ресурсы. Подходящий VPS-сервер позволяет разнести роли controller и broker по отдельным узлам, а не запускать все процессы на одной машине.
Подготовка VDS и сети
Команды установки и управления службами требуют административных прав. Такая схема подходит для VDS или выделенных серверов и не может быть развёрнута на обычном виртуальном хостинге.
Все узлы должны разрешать внутренние имена в приватные адреса. Предпочтителен внутренний DNS. Для небольшой тестовой инфраструктуры допустим одинаковый файл /etc/hosts на всех шести машинах:
10.20.0.11 controller-1.internal controller-1
10.20.0.12 controller-2.internal controller-2
10.20.0.13 controller-3.internal controller-3
10.20.0.21 broker-1.internal broker-1
10.20.0.22 broker-2.internal broker-2
10.20.0.23 broker-3.internal broker-3
Проверьте разрешение имён с каждого узла:
getent hosts controller-1.internal
getent hosts controller-2.internal
getent hosts controller-3.internal
getent hosts broker-1.internal
Не открывайте незащищённые Kafka listeners в публичный интернет. В примере используется PLAINTEXT, поэтому порты должны быть доступны только через доверенную приватную сеть. Контроллеры принимают TCP-подключения на 9093 от контроллеров и брокеров. Брокерам нужен TCP-порт 9092 для связи друг с другом и с разрешёнными клиентами.
Перед изменением правил выясните, какой файрвол уже работает. Не устанавливайте UFW поверх firewalld или собственного набора nftables:
sudo ufw status
sudo firewall-cmd --state
sudo nft list ruleset
Добавьте разрешения для приватных подсетей средствами действующего файрвола, не удаляя правила SSH. Например, в UFW правило для controller выглядит так:
sudo ufw allow from 10.20.0.0/24 to any port 9093 proto tcp
На broker аналогично разрешите 9092/tcp от подсети брокеров и клиентских приложений. Если приватная сеть работает по IPv6, добавьте правила для её IPv6-префикса и используйте соответствующие адреса в listeners. Kafka работает поверх TCP; открывать UDP-порты не требуется.
После применения правил проверьте маршрутизацию и доступность адресов. До запуска Kafka тест TCP-порта закономерно завершится ошибкой, но узлы должны разрешать имена и достигать друг друга по сети. Не отключайте действующие правила файрвола ради диагностики: это может открыть SSH или сервисы для нежелательных подключений.
Установка Java и Kafka
На всех узлах нужна одна и та же версия Kafka и поддерживаемая ею Java. Для Kafka 4.x практичным выбором является Java 21 LTS. Используйте актуальный пакет OpenJDK из репозитория дистрибутива.
Debian и Ubuntu
sudo apt update
sudo apt install openjdk-21-jre-headless
RHEL-совместимые системы
sudo dnf install java-21-openjdk-headless
Arch Linux
sudo pacman -Syu jre21-openjdk-headless
В Arch Linux эта команда обновляет систему. На уже используемом сервере заранее проверьте список пакетов и запланируйте перезагрузку, если обновились ядро или критические библиотеки.
Убедитесь, что Java доступна:
java -version
Загрузите официальный бинарный дистрибутив выбранной версии Kafka на каждый сервер, распакуйте его в отдельный каталог внутри /opt и создайте постоянную ссылку /opt/kafka. В командах замените имя каталога на фактическое:
sudo tar -xzf /tmp/kafka.tgz -C /opt
sudo ln -s /opt/kafka_2.13-x.y.z /opt/kafka
Если ссылка уже существует, сначала проверьте, куда она ведёт. Не заменяйте работающую версию Kafka без плана обновления и резервной копии конфигурации. Все шесть узлов этого кластера должны использовать одинаковую версию дистрибутива, иначе диагностика совместимости и последующее обслуживание усложнятся.
Создадим системного пользователя. Путь к nologin отличается между дистрибутивами, поэтому определим его автоматически:
NOLOGIN=$(command -v nologin)
sudo useradd --system --home-dir /var/lib/kafka --shell "$NOLOGIN" kafka
Если пользователь уже есть, useradd сообщит об этом — повторно создавать его не нужно. Подготовьте конфигурацию и хранилища:
sudo install -d -o root -g kafka -m 0750 /etc/kafka
sudo install -d -o kafka -g kafka -m 0750 /var/lib/kafka
sudo install -d -o kafka -g kafka -m 0750 /var/lib/kafka/metadata
Только на брокерах дополнительно создайте каталог данных:
sudo install -d -o kafka -g kafka -m 0750 /var/lib/kafka/data
На практике журналы данных брокеров часто размещают на отдельном виртуальном диске. В таком случае смонтируйте его заранее, внесите файловую систему в /etc/fstab по UUID и убедитесь, что каталог смонтирован до форматирования Kafka. После перезагрузки ещё раз проверьте точку монтирования: форматирование или запуск Kafka на несмонтированном каталоге может заполнить корневую файловую систему.
Настройка controller
На controller-1 создайте /etc/kafka/server.properties со следующим содержимым:
process.roles=controller
node.id=1001
controller.quorum.bootstrap.servers=controller-1.internal:9093,controller-2.internal:9093,controller-3.internal:9093
controller.listener.names=CONTROLLER
listeners=CONTROLLER://10.20.0.11:9093
listener.security.protocol.map=CONTROLLER:PLAINTEXT,INTERNAL:PLAINTEXT
inter.broker.listener.name=INTERNAL
metadata.log.dir=/var/lib/kafka/metadata
На остальных контроллерах используйте тот же файл, изменив только локальные значения:
- для
controller-2:node.id=1002иlisteners=CONTROLLER://10.20.0.12:9093; - для
controller-3:node.id=1003иlisteners=CONTROLLER://10.20.0.13:9093.
Список controller.quorum.bootstrap.servers должен быть одинаковым на controller и broker. Он помогает каждому процессу найти metadata quorum. Параметр process.roles=controller принципиален: не добавляйте к нему роль broker, иначе роли перестанут быть изолированными.
node.id должен быть уникальным в пределах всего кластера, а не только среди контроллеров. Нельзя назначить брокеру идентификатор, уже занятый controller. Также не меняйте node.id на уже отформатированном узле: он должен соответствовать идентификатору, записанному в локальном состоянии KRaft.
Закройте чтение конфигурации для посторонних пользователей:
sudo chown root:kafka /etc/kafka/server.properties
sudo chmod 0640 /etc/kafka/server.properties
Настройка broker
На broker-1 файл /etc/kafka/server.properties будет выглядеть так:
process.roles=broker
node.id=2001
controller.quorum.bootstrap.servers=controller-1.internal:9093,controller-2.internal:9093,controller-3.internal:9093
controller.listener.names=CONTROLLER
listeners=INTERNAL://10.20.0.21:9092
advertised.listeners=INTERNAL://broker-1.internal:9092
listener.security.protocol.map=CONTROLLER:PLAINTEXT,INTERNAL:PLAINTEXT
inter.broker.listener.name=INTERNAL
log.dirs=/var/lib/kafka/data
metadata.log.dir=/var/lib/kafka/metadata
num.partitions=3
default.replication.factor=3
min.insync.replicas=2
offsets.topic.replication.factor=3
transaction.state.log.replication.factor=3
transaction.state.log.min.isr=2
Для broker-2 измените node.id, адрес listener и advertised hostname:
node.id=2002
listeners=INTERNAL://10.20.0.22:9092
advertised.listeners=INTERNAL://broker-2.internal:9092
Для broker-3 используйте:
node.id=2003
listeners=INTERNAL://10.20.0.23:9092
advertised.listeners=INTERNAL://broker-3.internal:9092
listeners определяет локальный адрес, на котором процесс принимает соединения. advertised.listeners содержит адрес, возвращаемый клиентам в метаданных. Поэтому указанное там имя должно разрешаться не только на самом брокере, но и на всех клиентах Kafka.
Коэффициент репликации 3 предполагает три broker. Параметр min.insync.replicas=2 позволяет настроенным на подтверждение acks=all производителям продолжать запись при недоступности одного broker, но не при потере двух. Если брокеров меньше трёх, приведённые значения применять нельзя. Если позднее число broker изменится, параметры репликации и схему размещения данных нужно пересматривать отдельно.
Установите права на конфигурацию:
sudo chown root:kafka /etc/kafka/server.properties
sudo chmod 0640 /etc/kafka/server.properties
Создание cluster ID и идентификаторов каталогов controller
Новый KRaft-кластер должен получить один общий cluster ID. Его создают один раз и используют при форматировании всех controller и broker:
/opt/kafka/bin/kafka-storage.sh random-uuid
Сохраните результат в журнале развёртывания или защищённой системе управления конфигурациями. Для примера далее обозначим его как YOUR_CLUSTER_ID. Не публикуйте значение в открытых тикетах вместе с адресами инфраструктуры и не генерируйте новый ID для отдельных узлов.
Теперь создайте отдельный directory ID для каждого controller. Выполните команду три раза:
/opt/kafka/bin/kafka-storage.sh random-uuid
/opt/kafka/bin/kafka-storage.sh random-uuid
/opt/kafka/bin/kafka-storage.sh random-uuid
Предположим, получены значения CONTROLLER_1_DIRECTORY_ID, CONTROLLER_2_DIRECTORY_ID и CONTROLLER_3_DIRECTORY_ID. Из них формируется единый список первоначальных участников quorum:
1001@controller-1.internal:9093:CONTROLLER_1_DIRECTORY_ID,1002@controller-2.internal:9093:CONTROLLER_2_DIRECTORY_ID,1003@controller-3.internal:9093:CONTROLLER_3_DIRECTORY_ID
В каждой записи до первого символа @ указан node.id, затем адрес controller listener, а после последнего двоеточия — directory ID. Этот список вместе с cluster ID должен быть совершенно одинаковым при форматировании всех трёх контроллеров. Сопоставьте каждую запись с конфигурацией конкретного узла до запуска команды форматирования.
Форматирование хранилищ
Форматирование записывает служебные данные KRaft в указанные каталоги. Не запускайте его для рабочего узла повторно и не очищайте
metadata.log.dirв попытке устранить ошибку запуска. Сначала проверьте пути, права и cluster ID.
Так как инструкция относится к свежему кластеру, каталоги должны быть пустыми. Проверьте это на каждом узле:
sudo find /var/lib/kafka/metadata -mindepth 1 -maxdepth 2 -ls
На broker проверьте и каталог данных:
sudo find /var/lib/kafka/data -mindepth 1 -maxdepth 2 -ls
Если там есть неизвестные файлы, остановитесь и выясните их происхождение. Не удаляйте содержимое автоматически. До форматирования также проверьте владельца каталогов и доступность точки монтирования, если данные расположены на отдельном диске.
Форматирование controller
На каждом из трёх controller задайте одинаковые значения переменных, подставив реальные идентификаторы:
CLUSTER_ID='YOUR_CLUSTER_ID'
INITIAL_CONTROLLERS='1001@controller-1.internal:9093:CONTROLLER_1_DIRECTORY_ID,1002@controller-2.internal:9093:CONTROLLER_2_DIRECTORY_ID,1003@controller-3.internal:9093:CONTROLLER_3_DIRECTORY_ID'
Затем выполните:
sudo -u kafka /opt/kafka/bin/kafka-storage.sh format --cluster-id "$CLUSTER_ID" --initial-controllers "$INITIAL_CONTROLLERS" --config /etc/kafka/server.properties
Повторите команду на всех трёх controller. Она создаст локальный meta.properties и начальное состояние quorum. Не генерируйте новые значения для второго и третьего узла: cluster ID и полный список первоначальных controller должны остаться теми же.
После форматирования можно проверить локальный файл:
sudo -u kafka cat /var/lib/kafka/metadata/meta.properties
Cluster ID должен совпадать на всех узлах, а node.id — соответствовать конфигурации конкретного controller.
Форматирование broker
На каждом broker задайте тот же cluster ID:
CLUSTER_ID='YOUR_CLUSTER_ID'
Отформатируйте хранилище с флагом --no-initial-controllers:
sudo -u kafka /opt/kafka/bin/kafka-storage.sh format --cluster-id "$CLUSTER_ID" --config /etc/kafka/server.properties --no-initial-controllers
Этот флаг указывает, что broker присоединяется к уже определённому controller quorum и не создаёт его первоначальный состав. Выполните команду на всех трёх broker. Не запускайте broker до того, как будут отформатированы и запущены controller.
Создание службы systemd
На каждом узле создайте файл /etc/systemd/system/kafka.service:
[Unit]
Description=Apache Kafka Server
Wants=network-online.target
After=network-online.target
[Service]
Type=simple
User=kafka
Group=kafka
EnvironmentFile=-/etc/kafka/kafka-env
ExecStart=/opt/kafka/bin/kafka-server-start.sh /etc/kafka/server.properties
Restart=on-failure
RestartSec=10
TimeoutStopSec=180
LimitNOFILE=100000
SuccessExitStatus=143
[Install]
WantedBy=multi-user.target
Параметры heap удобно хранить отдельно. Начальный размер выбирайте по объёму памяти VDS и наблюдаемой нагрузке, оставляя память операционной системе и файловому кешу. Например:
KAFKA_HEAP_OPTS="-Xms1G -Xmx1G"
Сохраните строку в /etc/kafka/kafka-env, затем ограничьте доступ:
sudo chown root:kafka /etc/kafka/kafka-env
sudo chmod 0640 /etc/kafka/kafka-env
sudo systemctl daemon-reload
sudo systemctl enable kafka
Одинаковый unit можно использовать для обеих ролей: фактический режим процесса определяется файлом server.properties. Перед первым запуском проверьте, что unit ссылается на существующий путь /opt/kafka, а пользователь kafka имеет права записи в предназначенные для него каталоги.
Запуск metadata quorum
Сначала запустите все три controller, не включая broker:
sudo systemctl start kafka
sudo systemctl status kafka --no-pager
Команды выполняются на каждом controller. Если служба не перешла в состояние active (running), изучите журнал:
sudo journalctl -u kafka -n 100 --no-pager
Проверьте, что controller listener действительно открыт на каждом узле:
sudo ss -lntp | grep ':9093'
С другого controller или broker проверьте TCP-соединение:
nc -vz controller-1.internal 9093
nc -vz controller-2.internal 9093
nc -vz controller-3.internal 9093
Если nc не установлен, можно использовать другое доступное средство проверки TCP. Ошибка Connection refused обычно означает, что процесс не слушает порт. Тайм-аут чаще указывает на файрвол или сетевую маршрутизацию. Не переходите к запуску broker, пока controller не сформируют работоспособный quorum.
Проверка metadata quorum через CLI
Когда controller запущены, запросите состояние quorum с любого узла, где установлен тот же дистрибутив Kafka:
/opt/kafka/bin/kafka-metadata-quorum.sh --bootstrap-controller controller-1.internal:9093 describe --status
В выводе важны следующие поля:
ClusterIdдолжен совпадать с созданным идентификатором;LeaderIdдолжен содержать ID одного из controller, а не быть пустым;CurrentVotersдолжен включать три controller с ID1001,1002и1003;HighWatermarkпоказывает подтверждённую позицию metadata log;MaxFollowerLagпосле стабилизации кластера должен оставаться небольшим и не расти постоянно.
Более подробное состояние репликации выводится командой:
/opt/kafka/bin/kafka-metadata-quorum.sh --bootstrap-controller controller-1.internal:9093 describe --replication
Убедитесь, что followers догнали leader. Сразу после запуска возможен кратковременный lag, но он должен сокращаться. Если команда не подключается к controller, сначала проверьте DNS, доступность 9093/tcp и состояние службы на выбранном bootstrap-узле.
Запуск broker
После образования quorum запустите Kafka на каждом broker:
sudo systemctl start kafka
sudo systemctl status kafka --no-pager
Проверьте локальный listener:
sudo ss -lntp | grep ':9092'
Затем запросите версии API у broker:
/opt/kafka/bin/kafka-broker-api-versions.sh --bootstrap-server broker-1.internal:9092,broker-2.internal:9092,broker-3.internal:9092
Успешный вывод по всем трём адресам подтверждает, что broker запущены, зарегистрировались у controller и доступны клиенту. Проверку выполняйте с узла, который способен разрешить имена из advertised.listeners; иначе успешный локальный запуск не гарантирует доступность кластера приложениям.
Функциональная проверка кластера
Создадим тестовый топик с тремя разделами и коэффициентом репликации 3:
/opt/kafka/bin/kafka-topics.sh --bootstrap-server broker-1.internal:9092 --create --topic kraft-smoke --partitions 3 --replication-factor 3
Посмотрите распределение разделов:
/opt/kafka/bin/kafka-topics.sh --bootstrap-server broker-1.internal:9092 --describe --topic kraft-smoke
Для каждого раздела должны быть указаны три replicas. Список ISR после стабилизации также должен включать три broker ID. Если ISR неполон, проверьте журналы broker, доступность advertised.listeners и соединения по 9092/tcp.
Отправьте сообщение:
printf 'kraft-ok\n' | /opt/kafka/bin/kafka-console-producer.sh --bootstrap-server broker-1.internal:9092 --topic kraft-smoke
Прочитайте его через другой broker:
/opt/kafka/bin/kafka-console-consumer.sh --bootstrap-server broker-2.internal:9092 --topic kraft-smoke --from-beginning --max-messages 1
Ожидаемый результат — строка kraft-ok. После проверки тестовый топик можно удалить:
/opt/kafka/bin/kafka-topics.sh --bootstrap-server broker-1.internal:9092 --delete --topic kraft-smoke
Удаление топика выполняйте только после завершения проверки. Не используйте имя тестового топика, совпадающее с именем рабочего, и не выполняйте команду удаления в окружении, где bootstrap-адреса ведут в другой кластер.
Типичные ошибки при развёртывании
Quorum не выбирает leader
Чаще всего controller не видят друг друга из-за DNS, маршрутизации или файрвола. Проверьте 9093/tcp во всех направлениях, содержимое controller.quorum.bootstrap.servers и локальный адрес в listeners. Адрес listener должен существовать на конкретной VDS.
Другие причины — повторяющиеся node.id, разные cluster ID или неодинаковый список --initial-controllers при форматировании. Сравнивайте фактические конфигурации и значения, использованные при форматировании, а не только шаблоны файлов.
Ошибка несоответствия cluster ID
Она означает, что каталог уже отформатирован для другого кластера либо узел форматировали с ошибочным значением. Не исправляйте это удалением одного meta.properties: в каталоге остаётся связанное состояние.
Если кластер уже содержит полезные данные, остановитесь и восстановите правильную конфигурацию или хранилище из предусмотренной процедуры восстановления. Для полностью нового стенда без данных допустима только полная повторная инициализация: остановить все шесть процессов, сохранить конфигурации и журналы, точно проверить пути, очистить исключительно выделенные каталоги Kafka на всех узлах, создать новый комплект ID и заново отформатировать весь кластер.
Broker запускается, но клиенты теряют соединение
Проверьте advertised.listeners. Начальное подключение к bootstrap server может пройти успешно, после чего клиент получит адреса остальных broker. Если эти имена недоступны клиенту, дальнейшая работа завершится сетевой ошибкой. Исправлять нужно адреса, DNS или маршрутизацию, а не только список bootstrap-серверов у клиента.
Kafka не может записать данные
Проверьте владельца каталогов, наличие свободного места и факт монтирования диска:
sudo -u kafka test -w /var/lib/kafka/metadata && echo writable
df -h /var/lib/kafka/metadata
findmnt /var/lib/kafka/metadata
На broker повторите проверку для /var/lib/kafka/data. Если отдельный диск не смонтировался, Kafka может записать данные в каталог на корневой файловой системе. Это способно заполнить системный раздел, поэтому состояние mount point следует контролировать до запуска службы. Для организации наблюдения за дисковым пространством и доступностью узлов пригодится материал о мониторинге VDS с Prometheus и Node Exporter.
Служба постоянно перезапускается
Временно остановите цикл перезапуска и изучите первую ошибку, а не только последние сообщения:
sudo systemctl stop kafka
sudo journalctl -u kafka --since '15 minutes ago' --no-pager
После исправления конфигурации запустите сервис вручную через systemd. Не запускайте второй экземпляр kafka-server-start.sh параллельно службе: он попытается использовать те же порты и каталоги.
Что проверить перед передачей кластера в работу
- На каждом узле установлены одинаковые версии Kafka и совместимой Java.
- Все шесть
node.idуникальны. - На controller задана только роль
controller, на broker — толькоbroker. - Cluster ID одинаков во всех файлах
meta.properties. - Metadata quorum содержит три voters и имеет leader.
- Followers metadata log не демонстрируют постоянно растущий lag.
- Три broker доступны через адреса из
advertised.listeners. - Порты
9092и9093не открыты в публичную сеть при использованииPLAINTEXT. - Тестовый топик получает три replicas и полный ISR.
- Сообщение успешно записывается через один broker и читается через другой.
- Диски с
metadata.log.dirиlog.dirsимеют мониторинг свободного места. - Конфигурации, cluster ID и соответствие controller directory ID сохранены в эксплуатационной документации.
В результате controller независимо поддерживают KRaft metadata quorum, а broker можно обслуживать и масштабировать отдельно. Такая схема сложнее совмещённых ролей, зато нагрузка и перезапуск broker не затрагивают процесс controller на той же VDS, а состав каждого слоя можно изменять по собственной процедуре.


