Apache Kafka на сервере: частые ошибки и решения
Kafka на локальной машине с одним брокером и дефолтным конфигом поднимается за пять минут — а на боевом сервере та же Kafka вдруг отказывается стартовать после перезагрузки, диск забивается сегментами лога за пару дней, а клиенты снаружи вообще не могут подключиться к брокеру. Разберём шесть типовых причин отказов Kafka на сервере — от несовпадения cluster ID до неверных advertised listeners — с конкретными командами диагностики и рабочими исправлениями.
Содержание
Обсудить статью, задать вопрос или начать новую тему
Есть вопрос по этой статье, идея для обсуждения или просто хотите поделиться опытом? Сообщество MAATRIX ждёт. Для общения, пожалуйста, зарегистрируйтесь в нашем личном кабинете.
Перейти в сообщество →Брокер не стартует: cluster ID и meta.properties
Начиная с версии 3.x Kafka официально поддерживает режим KRaft без ZooKeeper, а с 4.0 ZooKeeper убран из дистрибутива вовсе — брокер сам выполняет роли контроллера через consensus-протокол Raft. Отсюда и самая частая ошибка первого запуска: попытка стартовать брокер с пустым или неинициализированным каталогом данных. Kafka в режиме KRaft требует явного форматирования каталога логов перед первым стартом:
KAFKA_CLUSTER_ID="$(kafka-storage.sh random-uuid)"
kafka-storage.sh format -t $KAFKA_CLUSTER_ID -c /opt/kafka/config/server.properties
Если пропустить этот шаг, брокер упадёт с ошибкой вида Log directory ... is not formatted. Вторая частая ситуация — брокер уже работал, но при восстановлении из бэкапа, клонировании диска на другой сервер или ручном копировании log.dirs вы получаете InconsistentClusterIdException: The Cluster ID doesn't match stored clusterId. Kafka пишет cluster.id в файл meta.properties внутри каждого каталога логов и жёстко сверяет его при каждом старте — это защита от случайного присоединения диска не от того кластера:
cat /var/lib/kafka/data/meta.properties
Если это действительно новый изолированный кластер, а не восстановление существующего — просто отформатируйте каталог заново командой выше с новым KAFKA_CLUSTER_ID. Если это должен быть тот же кластер — восстановите правильный meta.properties вместо генерации нового, иначе брокер создаст себе новую идентичность и не сможет присоединиться к остальным нодам.
Диск заполняется: retention и сегменты лога
Kafka физически хранит сообщения на диске как последовательность сегментных файлов, и без явных настроек retention они будут копиться бесконечно — Kafka не удаляет данные автоматически после того, как их прочитал последний консьюмер, удаление регулируется только политикой retention по времени или объёму:
log.retention.hours=168
log.retention.bytes=-1
log.segment.bytes=1073741824
log.retention.check.interval.ms=300000
Частая ошибка — оставить log.retention.bytes=-1 и полагаться только на retention.hours, а затем удивляться, что при всплеске трафика диск заполняется быстрее, чем истекает недельный retention. Для топиков с предсказуемо большим объёмом лучше явно задать оба лимита — сработает тот, что наступит раньше. Проверить реальное потребление диска по каждому топику:
du -sh /var/lib/kafka/data/*/
kafka-log-dirs.sh --bootstrap-server localhost:9092 --describe | jq .
Если брокер всё же упёрся в No space left on device, он останавливает запись в соответствующий каталог логов и может пометить его offline — новые партиции туда не разместятся, пока место не освободится, а старые сегменты руками удалять небезопасно (можно повредить индекс сегмента). Правильный порядок: временно ужесточить retention на самом «тяжёлом» топике через kafka-configs.sh, дать брокеру самому вычистить старые сегменты по расписанию, и только после этого вернуть retention к нормальным значениям. За общим состоянием диска на проде удобно следить постоянно, а не по факту падения — подробнее в статье про мониторинг диска на сервере.
Нужен сервер под эту задачу?
Разверните VPS MAATRIX за пару минут: NVMe, AMD EPYC, root-доступ, локации UK, США, Франция и РФ. Оплата картой РФ и по СБП.
Арендовать серверJVM heap и память: зачем Kafka нужен page cache
Kafka — JVM-приложение, но, в отличие от классических баз данных, она сознательно спроектирована так, чтобы полагаться не на кучу (heap), а на файловый кеш операционной системы. Брокер пишет и читает сегменты через обычные файловые операции ОС, и именно page cache, а не heap, отвечает за то, что чтение недавних сообщений происходит без обращения к диску. Отсюда практическое правило: heap Kafka-брокера почти никогда не должен быть больше 6-8 ГБ, даже на сервере с 32-64 ГБ RAM — оставшуюся память нужно оставить операционной системе под page cache:
export KAFKA_HEAP_OPTS="-Xms6g -Xmx6g"
export KAFKA_JVM_PERFORMANCE_OPTS="-XX:+UseG1GC -XX:MaxGCPauseMillis=20 -XX:InitiatingHeapOccupancyPercent=35"
Частая ошибка новичков — увеличивать heap «про запас» до половины RAM сервера, как советуют делать для Elasticsearch или других JVM-хранилищ. Для Kafka это контрпродуктивно: чем больше heap, тем меньше page cache, тем чаще брокер реально читает с диска вместо памяти, и тем сильнее просадка на GC-паузах при большой куче. Если процесс брокера внезапно исчезает без единой строчки в собственном логе — почти всегда это OOM killer ядра Linux, а не JVM:
sudo dmesg -T | grep -i "killed process"
sudo journalctl -u kafka --since "2 hours ago" | grep -i oom
Такое случается, когда heap задан правильно, но на том же сервере одновременно крутятся другие тяжёлые сервисы и суммарное потребление превышает физическую память. Дальше упирается уже не конфиг Kafka, а объём RAM конкретной машины — память можно нарастить под растущую нагрузку заранее, не разнося сервисы по разным хостам в аварийном режиме.
Under-replicated partitions и выпадение из ISR
Каждая партиция в Kafka имеет одного лидера и несколько реплик, а список реплик, реально успевающих за лидером, называется ISR (in-sync replicas). Партиция считается under-replicated, если хотя бы одна реплика из назначенных отстала и выпала из ISR — это главный сигнал нездоровья кластера, который стоит мониторить постоянно:
kafka-topics.sh --bootstrap-server localhost:9092 --describe --under-replicated-partitions
Типичные причины выпадения реплики из ISR — не сетевой сбой, а банальная нехватка ресурсов: медленный диск на одном из брокеров, который не успевает записывать сегменты с той же скоростью, что лидер; долгая GC-пауза из-за завышенного heap (см. предыдущий раздел); или сетевая задержка между брокерами, если они физически разнесены по разным дата-центрам без учёта латентности. Порог, после которого реплика считается отставшей, задаётся параметром:
replica.lag.time.max.ms=30000
Соблазн — просто увеличить это значение, чтобы «медленные» реплики не выпадали из ISR. Это лечит симптом, а не причину: реплика, которая реально не успевает, будет тормозить подтверждение записи для продюсеров с acks=all, просто позже. Правильный порядок — сначала найти узкое место конкретного брокера (диск, сеть, CPU, GC), а параметр replica.lag.time.max.ms трогать только если задержка объясняется законной сетевой топологией, а не деградацией железа. Общее состояние кластера удобно вынести на дашборд — Kafka отдаёт метрики через JMX, которые снимает jmx_exporter для Prometheus, а визуализация делается той же связкой, что и для других сервисов — см. Grafana и Prometheus на сервере.
Advertised listeners: клиенты не видят брокер снаружи
Едва ли не самая частая практическая проблема при первом развёртывании Kafka на VPS — брокер вроде бы работает, kafka-topics.sh --list с самого сервера отвечает нормально, но клиент с внешней машины не может ни подключиться, ни получить список партиций, хотя порт 9092 открыт файрволом. Причина — путаница между listeners (на каких адресах брокер слушает сокеты) и advertised.listeners (какой адрес брокер сообщает клиентам как «подключайтесь сюда» при первом handshake). Клиент сначала стучится на bootstrap-адрес, получает в ответ metadata с реальным advertised-адресом и уже к нему открывает соединения для чтения/записи — если advertised-адрес недоступен снаружи, всё ломается на втором шаге, а не на первом:
listeners=INTERNAL://0.0.0.0:9092,EXTERNAL://0.0.0.0:9093
advertised.listeners=INTERNAL://kafka-1.internal:9092,EXTERNAL://203.0.113.10:9093
listener.security.protocol.map=INTERNAL:PLAINTEXT,EXTERNAL:PLAINTEXT
inter.broker.listener.name=INTERNAL
Классическая ловушка в Docker: если оставить advertised.listeners со значением localhost или внутренним именем контейнера по умолчанию, клиент снаружи Docker-сети получает это имя в metadata и пытается резолвить его на своей машине — безуспешно. Отдельный листенер EXTERNAL с публичным IP или доменом сервера и своим портом решает проблему аккуратно, не смешивая межброкерный трафик с внешним трафиком клиентов. Если Kafka развёрнута через docker-compose, тот же принцип действует и там — общий подход к конфигурации портов и сети контейнеров разобран в статье про docker-compose для продакшена на сервере. Не забудьте и про файрвол — порты EXTERNAL-листенера должны быть явно открыты, а порт межброкерного трафика (INTERNAL) лучше вообще не выставлять наружу.
Producer и consumer: таймауты и растущий lag
Ошибка продюсера NotLeaderOrFollowerException или UnknownTopicOrPartitionException в момент ребалансировки партиций — это нормальная переходная ситуация, а не сбой: пока Kafka выбирает нового лидера партиции после падения брокера, короткое окно запросы будут отбиваться, и настроенные по умолчанию ретраи продюсера должны это окно пережить сами:
retries=2147483647
delivery.timeout.ms=120000
request.timeout.ms=30000
enable.idempotence=true
Хуже, если ошибки продюсера не разовые, а постоянные под нагрузкой — тогда стоит проверить request.timeout.ms относительно реальной задержки диска на брокерах (см. раздел про ISR) и включить enable.idempotence=true — она почти не стоит ничего по производительности, но защищает от дублей при повторных отправках на нестабильной сети.
Со стороны потребителей главная метрика здоровья — не «работает / не работает», а consumer lag: разница между последним записанным офсетом партиции и офсетом, до которого дочитала конкретная consumer group:
kafka-consumer-groups.sh --bootstrap-server localhost:9092 --describe --group my-group
Растущий лаг обычно означает одно из трёх: потребитель реально не успевает обрабатывать сообщения (нужно больше партиций плюс больше инстансов consumer group либо ускорить саму обработку), у него слишком короткий session.timeout.ms и он попадает в бесконечные ребалансировки, теряя время не на обработку, а на переприсоединение к группе, либо max.poll.interval.ms меньше реального времени обработки батча — тогда координатор считает потребителя мёртвым посреди работы и выгоняет его из группы. Штормы ребалансировки в логах брокера, повторяющиеся с интервалом в секунды — почти всегда сигнал разобраться с таймаутами консьюмера, а не увеличивать число партиций вслепую.
Нужен сервер под эту задачу?
Разверните VPS MAATRIX за пару минут: NVMe, AMD EPYC, root-доступ, локации UK, США, Франция и РФ. Оплата картой РФ и по СБП.
Арендовать серверНужны сами нейросети для контента?
Генерируйте изображения, видео и озвучку нейросетями на falapi.io — десятки моделей в одном окне. Оплата картой РФ и по СБП.
Частые вопросы
Kafka не стартует после переноса на новый сервер — что делать?
Проверьте meta.properties в каталоге логов — скорее всего, cluster.id не совпадает с тем, что брокер ожидает при старте. Если это действительно новый изолированный кластер, отформатируйте каталог заново через kafka-storage.sh format; если нужен тот же кластер — восстановите оригинальный meta.properties, а не создавайте новый.
Сколько RAM реально нужно под Kafka-брокер на проде?
Отправная точка — heap 6-8 ГБ независимо от общего объёма RAM сервера, а остальную память нужно оставлять под page cache операционной системы: именно он, а не heap, ускоряет чтение недавних сообщений. Для серьёзной нагрузки чаще не хватает не heap, а диска и его скорости.
Почему клиент с другого сервера не может подключиться к Kafka, хотя порт открыт?
Скорее всего дело не в файрволе, а в advertised.listeners — брокер сообщает клиенту неверный (внутренний или localhost) адрес в ответ на первый запрос metadata. Настройте отдельный листенер с публичным адресом для внешних клиентов.
Under-replicated partitions — это критично?
Да, это главный индикатор нездоровья кластера, и его стоит мониторить постоянно, а не разбирать по факту жалоб. Причина почти всегда в ресурсах конкретного брокера (диск, GC-паузы, сеть), а не в самой Kafka.
Нужно ли увеличивать число партиций топика при росте лага у консьюмеров?
Не всегда — сначала проверьте session.timeout.ms и max.poll.interval.ms: если лаг растёт из-за постоянных ребалансировок, а не из-за реальной нехватки параллелизма, увеличение числа партиций проблему не решит.
Обсудить статью, задать вопрос или начать новую тему
Есть вопрос по этой статье, идея для обсуждения или просто хотите поделиться опытом? Сообщество MAATRIX ждёт. Для общения, пожалуйста, зарегистрируйтесь в нашем личном кабинете.
Перейти в сообщество →