Перейти к содержимому

Kafka: проверка работоспособности кластера

Kafka-кластер в KRaft-режиме не прощает, когда о нём забывают до первого инцидента. Проверять здоровье нужно регулярно и короткими командами — без графиков и дашбордов, прямо из консоли. Ниже — набор команд, которые закрывают типовой чек-лист: процессы, кворум, лидеры партиций, ISR и быстрый «светофор» одним вызовом.

Проверка процессов KRaft

В режиме KRaft отдельного ZooKeeper нет, и роль controller совмещена с ролью broker или вынесена на отдельные ноды. Сначала убедимся, что JVM-процессы вообще живы, и посмотрим, в каком режиме стартовал каждый узел.

ps -ef | grep -E 'kafka.Kafka|QuorumControllerMain' | grep -v grep

Для управляемого кластера ожидаемо увидеть три процесса QuorumControllerMain на выделенных контроллерах и N процессов kafka.Kafka на брокерах. На смешанных нодах (combined mode) будет один процесс сразу с двумя ролями — это нормально для небольших инсталляций.

Дальше — журнал запуска и конфигурация:

grep -E 'kafka-process-roles|node.id|controller.quorum.voters|listeners=' config/kraft/server.properties
Примечание

Параметр process.roles принимает значения broker, controller или broker,controller. Если в нём пусто — кластер ещё в legacy-режиме с ZooKeeper, и команды ниже нужно адаптировать.

Проверяем, что все ноды договорились о кворуме:

/opt/kafka/bin/kafka-metadata-quorum.sh \
  --bootstrap-server localhost:9092 describe --status

В выводе ищем leaderId, votedLeaders и размер кворума 2/3 (для трёх контроллеров) или N/N для уже стабилизированного кластера.

Статусы брокеров и контроллера

Дальше — понять, кто из брокеров реально отвечает, а кто выпал из реестра. Используем kafka-broker-api-versions.sh: он возвращает поддерживаемые API-версии и попутно показывает, доходит ли TCP-соединение до брокера.

for h in kafka1 kafka3 kafka5; do
  echo "=== $h ==="
  /opt/kafka/bin/kafka-broker-api-versions.sh \
    --bootstrap-server $h:9092 2>&1 | head -n 3
done

Если узел недоступен — увидим таймаут. Это самый быстрый способ отличить «брокер висит в JVM» от «сеть режет».

Полный список зарегистрированных брокеров и их состояние:

/opt/kafka/bin/kafka-broker-api-versions.sh \
  --bootstrap-server kafka1:9092 | head -n 50

Команда покажет JSON-like листинг, но для табличного отчёта удобнее kafka-metadata-quorum.sh:

/opt/kafka/bin/kafka-metadata-quorum.sh \
  --bootstrap-server kafka1:9092 describe --replicas

Колонка LEADER показывает текущий лидер кворума, REPLICAS — все активные узлы. Контроллеры, не отвечающие на запрос, выпадут из списка.

Подсказка

Поле lastCaughtUpTime в выводе describe --status показывает, насколько контроллер отстал от лидера. Значение 0 или свежее now — здоров, отставание в минутах — повод смотреть GC и сетевые задержки.

Список топиков и partition leaders

Кластер может быть жив, но без лидеров партиций — продюсер не запишет, консьюмер не прочитает. Поэтому следующий шаг — топики и их лидеры.

Список всех топиков с количеством партиций и реплик:

/opt/kafka/bin/kafka-topics.sh \
  --bootstrap-server kafka1:9092 \
  --describe --exclude-internal

Таблица вывода содержит Leader, Replicas, Isr. Если Leader равен -1, значит, партиция не имеет активного лидера — это авария, продюсеры будут получать NotLeaderForPartitionException.

Чтобы получить только «плохие» партиции в одном пайпе:

/opt/kafka/bin/kafka-topics.sh \
  --bootstrap-server kafka1:9092 \
  --describe --exclude-internal \
  | awk '$5 == -1 || $5 == "none" {print}'

Если хочется видеть лидеров по конкретному топику:

/opt/kafka/bin/kafka-topics.sh \
  --bootstrap-server kafka1:9092 \
  --describe --topic orders.events

ISR и недоступные реплики

ISR (in-sync replicas) — это то, что определяет надёжность записи. Рекомендуется держать min.insync.replicas >= 2 для критичных топиков. Проверяем рассинхрон:

/opt/kafka/bin/kafka-topics.sh \
  --bootstrap-server kafka1:9092 \
  --describe --under-replicated-partitions

Команда вернёт партиции, у которых Isr меньше, чем Replicas. Пустой вывод — хорошо. Любая строка в списке — инцидент.

Полная картина по всем партициям с фильтром по проблемам:

/opt/kafka/bin/kafka-topics.sh \
  --bootstrap-server kafka1:9092 \
  --describe --exclude-internal \
  | awk '{
    replicas=$5; isr=$7;
    # Replicas и Isr — списки id через запятую, длина считается по запятым
    rep_n = split(replicas, a, ",");
    isr_n = split(isr, b, ",");
    if (rep_n != isr_n) print "UNSYNC:", $0;
    if (a[1] == "-1") print "NO_LEADER:", $0;
  }'
Предупреждение

kafka-topics.sh --describe для топика с тысячами партиций выводит много строк и нагружает контроллер. На проде запускайте с --partitions N или фильтруйте awk, иначе чек сам станет источником проблемы.

Дополнительно — состояние конкретной реплики на стороне брокера. Если подозреваем, что один из дисков отстал:

/opt/kafka/bin/kafka-log-dirs.sh \
  --bootstrap-server kafka1:9092 \
  --describe --broker-list 1,3,5 \
  | jq '.[] | .logDirs[] | {broker: .broker, dir: .dir, partitions: (.partitions | length)}'

В выводе смотрим поле partition.error — если оно непустое, реплика имеет проблемы (offline log dir, диск переполнен, fs в read-only).

Быстрая диагностика в одной команде

Для ежедневных обходов удобно собрать все проверки в одном скрипте с понятными exit-кодами. Ниже — минимальный «светофор» на bash.

#!/usr/bin/env bash
set -u
BOOTSTRAP="${BOOTSTRAP:-kafka1:9092}"
KAFKA_BIN="${KAFKA_BIN:-/opt/kafka/bin}"
fail=0

echo "== Quorum status =="
if ! $KAFKA_BIN/kafka-metadata-quorum.sh --bootstrap-server "$BOOTSTRAP" \
    describe --status 2>&1 | grep -q 'isLeader: true'; then
  echo "WARN: quorum leader not confirmed"; fail=1
fi

echo "== Brokers reachability =="
for h in $(echo "$BOOTSTRAP" | tr ',' ' '); do
  if ! timeout 5 bash -c "echo > /dev/tcp/${h%:*}/${h##*:}"; then
    echo "FAIL: $h unreachable"; fail=1
  fi
done

echo "== Under-replicated partitions =="
out=$($KAFKA_BIN/kafka-topics.sh --bootstrap-server "$BOOTSTRAP" \
      --describe --under-replicated-partitions 2>/dev/null)
if [[ -n "$out" ]]; then
  echo "$out"; fail=1
else
  echo "OK: 0 under-replicated partitions"
fi

echo "== Partitions without leader =="
out=$($KAFKA_BIN/kafka-topics.sh --bootstrap-server "$BOOTSTRAP" \
      --describe --exclude-internal 2>/dev/null \
      | awk '$5 == -1 {print}')
if [[ -n "$out" ]]; then
  echo "$out"; fail=1
else
  echo "OK: every partition has a leader"
fi

exit $fail

Сохраняем как kafka-health.sh, делаем исполняемым и заворачиваем в cron или systemd timer раз в 60 секунд:

chmod +x kafka-health.sh
*/1 * * * * /usr/local/bin/kafka-health.sh \
  >> /var/log/kafka-health.log 2>&1
Подсказка

Для алертинга замените блок echo на logger -p local0.err и настройте rsyslog в SIEM. Так Prometheus node_exporter не нужен — события летят в общий канал.

Типичные ошибки при первом прогоне и как их трактовать:

СимптомВероятная причинаЧто делать
Connection to node -1 could not be establishedБрокер не зарегистрирован в кластере, но процесс живПроверить advertised.listeners, node.id, сетевой ACL
isLeader: false у всех контроллеровПотерян кворум, контроллеров < __.min.insync.replicasСмотреть controller.quorum.voters и состояние дисков на контроллерах
--under-replicated-partitions показывает записи дольше 5 минутБрокер отстаёт, диск медленный или GC-паузыСнять jstack, проверить iostat -x и метрики LogFlushRate
kafka-log-dirs.sh возвращает LogDirOfflineДиск переполнен или вышел из строяОсвободить место, проверить ФС на read-only, рестартовать брокер
Пустой вывод --describe при работающем кластереПередан неверный bootstrap-адресСверить advertised.listeners и DNS

Пять команд выше покрывают 90% оперативных вопросов «а кластер живой?». Если все зелёные — копать глубже не нужно; если что-то красное — kafka-log-dirs.sh и kafka-metadata-quorum.sh describe --status покажут, в какую сторону двигаться.