Apache Kafka: архитектура, ключевые понятия и работа с кластером | AdminWiki

Apache Kafka: архитектура, ключевые понятия и работа с кластером

18 августа 2026 16 мин. чтения
Содержание статьи

Введение в Apache Kafka: зачем нужна распределённая платформа потоковой передачи данных

Apache Kafka - это распределённая платформа потоковой передачи данных. Она принимает, хранит и отдаёт потоки записей в реальном времени. Kafka решает задачу надёжной доставки больших объёмов данных между системами: от источников к потребителям, с минимальной задержкой. Гарантия сохранности зависит от replication factor, состояния ISR и настроек производителя, поэтому для production важны не только базовые понятия, но и правильная проверка кластера.

Ключевые преимущества Kafka: высокая пропускная способность, горизонтальная масштабируемость и отказоустойчивость. Один кластер обрабатывает миллионы сообщений в секунду. Данные хранятся на диске, а не в оперативной памяти, поэтому Kafka не теряет их при перезапуске при корректно настроенной репликации и может удерживать большие объёмы без раздувания RAM.

Типичные сценарии использования: сбор логов с сотен серверов, потоковая обработка событий, интеграция микросервисов, аналитика в реальном времени. Если вам нужно сравнить Kafka с RabbitMQ или NATS для высоких нагрузок, на сайте есть отдельное практическое сравнение брокеров сообщений.

Ниже последовательно разобраны архитектура Kafka, работа KRaft и controller quorum, хранение данных, настройка Kafka-кластера из нескольких брокеров, CLI Kafka, репликация, проверка здоровья и диагностика типовых инцидентов.

Архитектура Kafka: брокеры, ZooKeeper и KRaft

Кластер Kafka состоит из брокеров - серверов, которые хранят данные и обрабатывают запросы клиентов. Брокеры объединяются в кластер и распределяют между собой нагрузку. Для координации метаданных используются два механизма: ZooKeeper в старых версиях и KRaft в новых.

Брокеры: сердце кластера Kafka

Брокер - это процесс Kafka, запущенный на сервере. Каждый брокер хранит одну или несколько партиций топиков. Он принимает записи от производителей, сохраняет их на диск и отдаёт потребителям. Брокеры обмениваются данными между собой для репликации: лидер партиции передаёт новые записи своим последователям.

Каждый брокер имеет уникальный числовой идентификатор broker.id. Этот параметр задаётся в конфигурационном файле server.properties. В кластере из трёх брокеров идентификаторы могут быть 0, 1 и 2. Брокеры периодически обмениваются heartbeat-сообщениями, чтобы контроллер знал, какие узлы живы.

ZooKeeper и KRaft: управление метаданными кластера

ZooKeeper - это внешний сервис координации, который использовался в Kafka до перехода на KRaft. Он хранил метаданные: список топиков, расположение партиций, информацию о лидерах и ISR. ZooKeeper также отвечал за выбор контроллера - брокера, который управляет репликацией и переизбранием лидеров.

KRaft (Kafka Raft) - встроенный протокол управления кластером, который устраняет зависимость от ZooKeeper. В режиме KRaft метаданные хранятся внутри самого кластера Kafka с использованием алгоритма консенсуса Raft. Это упрощает развертывание: не нужно поднимать отдельный ZooKeeper-ансамбль. Начиная с Kafka 3.3 режим KRaft считается production-ready, а в версии 4.0 ZooKeeper полностью удалён. Для новых инсталляций рекомендуется сразу использовать KRaft.

В KRaft controller quorum состоит из controller-узлов, которые участвуют в консенсусе Raft и хранят журнал метаданных. Брокеры отвечают за хранение партиций и клиентские запросы, а контроллеры управляют состоянием кластера, назначением лидеров и изменениями метаданных. В небольшом production-кластере роли можно совместить на трёх узлах: каждый узел работает как broker и controller, а quorum из трёх контроллеров переживает отказ одного узла.

При росте нагрузки роли обычно разделяют: используют не менее трёх controller-узлов и отдельные broker-узлы. Это уменьшает влияние дисковых операций, сетевых запросов и восстановления партиций на работу controller quorum. Для кластера из нечётного числа контроллеров проще обеспечить кворум: три узла выдерживают отказ одного, пять - отказ двух.

При запуске KRaft особенно важны параметры process.roles, node.id, controller.quorum.voters, controller.listener.names, listeners, advertised.listeners и inter.broker.listener.name. Все узлы должны использовать один идентификатор кластера и одинаковый список controller quorum, но иметь собственный node.id, адреса и каталог данных.

Ключевые понятия Kafka: топики, партиции, смещения и группы потребителей

Четыре концепции лежат в основе работы с Kafka: топики, партиции, смещения и группы потребителей. Без чёткого понимания этих понятий невозможно проектировать продюсеров и консьюмеров.

Топики и партиции: организация потоков данных

Топик - это логическая категория, в которую производители публикуют сообщения. Например, топик orders может содержать все события о заказах. Топик делится на партиции - физические подразделения, которые обеспечивают параллельную обработку.

Количество партиций задаётся при создании топика и напрямую влияет на пропускную способность: чем больше партиций, тем больше потребителей могут читать топик параллельно. Партиции одного топика могут находиться на разных брокерах. Это позволяет распределять нагрузку по кластеру. Увеличить число партиций можно позже, но уменьшить нельзя.

Планирование партиций и replication factor

Число партиций выбирают по результатам измерений, а не только по размеру текущего трафика. Практический ориентир - взять большее значение из требуемого числа параллельных потребителей и расчёта требуемая пропускная способность / пропускная способность одной партиции, затем добавить запас на рост нагрузки. Фактическую пропускную способность одной партиции нужно определить нагрузочным тестом для конкретного размера сообщений, ключей, дисков и сети.

Выводы только по throughput недостаточны. Слишком большое число партиций увеличивает объём метаданных, количество открытых файлов, время ребалансировки consumer group и восстановления после отказа брокера. При планировании также учитывайте пиковую нагрузку, retention, допустимое время восстановления и максимальное число потребителей в группе.

Replication factor задаёт число копий каждой партиции. Для небольшого production-кластера обычно используют replication.factor=3 на трёх и более брокерах и размещают реплики в разных failure domain, если это поддерживает инфраструктура. Значение 1 не обеспечивает отказоустойчивость, а значение 2 оставляет меньше запаса при одновременных проблемах. replication factor не может быть больше числа доступных брокеров, а увеличение числа копий повышает требования к диску и межброкерной сети.

Смещения (offset): позиция сообщения в партиции

Смещение - это уникальный последовательный идентификатор сообщения внутри партиции. Первое сообщение получает offset 0, второе - 1, и так далее. Смещение не используется для удаления сообщений: данные удаляются по времени хранения или размеру, а не по факту прочтения.

Потребитель сохраняет своё текущее смещение, чтобы после перезапуска продолжить чтение с нужного места. По умолчанию Kafka хранит смещения в специальном внутреннем топике __consumer_offsets. Это позволяет консьюмеру восстановить позицию даже после сбоя.

Группы потребителей (consumer groups): параллельная обработка сообщений

Группа потребителей - это набор консьюмеров, которые совместно читают один или несколько топиков. Kafka назначает каждую партицию только одному потребителю в группе. Если в группе три потребителя и топик имеет шесть партиций, каждый потребитель получит по две партиции.

Если потребителей больше, чем партиций, часть из них будет простаивать. При добавлении или удалении потребителя из группы Kafka запускает ребалансировку: перераспределяет партиции между активными участниками. Во время ребалансировки чтение приостанавливается, поэтому частые изменения состава группы нежелательны.

Хранение данных в Kafka: как достигается высокая производительность

Kafka хранит данные на диске в виде логов. Каждая партиция - это append-only лог: записи добавляются только в конец. Последовательная запись на диск работает быстрее случайной, поэтому Kafka достигает высокой пропускной способности даже на обычных HDD.

Для потоков с жёсткими требованиями к задержке обычно выбирают SSD или NVMe. HDD могут обеспечить приемлемый throughput для последовательных операций, но задержка, время восстановления и поведение при параллельной работе с большим числом партиций будут менее предсказуемыми.

Логи и сегменты: структура хранения партиции

Партиция физически состоит из нескольких сегментов - файлов фиксированного размера. По умолчанию размер сегмента - 1 ГБ. Когда сегмент заполняется, Kafka создаёт новый. Старые сегменты удаляются или архивируются согласно политике хранения, которая задаётся параметрами log.retention.hours или log.retention.bytes.

Каждый сегмент - это два файла: .log с данными и .index с индексами смещений. Индекс позволяет быстро находить сообщение по offset без полного сканирования сегмента.

Механизмы высокой пропускной способности: zero-copy и page cache

Zero-copy - это техника передачи данных из файловой системы в сеть без копирования в пользовательское пространство. При чтении Kafka использует системный вызов sendfile(), который передаёт данные напрямую из page cache в сокет. Это снижает нагрузку на CPU и ускоряет доставку сообщений потребителям.

Page cache - это кэш операционной системы, который хранит недавно прочитанные или записанные страницы диска в оперативной памяти. Kafka активно использует page cache: при записи данные попадают в кэш и асинхронно сбрасываются на диск. При чтении, если данные ещё в кэше, Kafka отдаёт их без обращения к диску. Это даёт скорость, близкую к in-memory системам, при сохранении надёжности дискового хранения. Фактическая защита от отказа узла обеспечивается прежде всего репликацией и корректными настройками подтверждения записи.

Развертывание и настройка кластера Kafka

Для production-кластера минимально рекомендуются три брокера. Это обеспечивает отказоустойчивость при выходе из строя одного узла. Ниже - пошаговая инструкция по установке и настройке.

Установка и запуск Kafka: пошаговая инструкция

Для работы Kafka требуется Java 11 или новее. Проверьте версию:

java -version

Скачайте официальный дистрибутив Kafka с сайта Apache и распакуйте архив:

tar -xzf kafka_2.13-3.7.0.tgz
cd kafka_2.13-3.7.0

Для режима KRaft сгенерируйте идентификатор кластера и отформатируйте хранилище:

bin/kafka-storage.sh random-uuid
bin/kafka-storage.sh format -t <uuid> -c config/kraft/server.properties

Команда random-uuid выводит идентификатор кластера. Сохраните его и используйте один и тот же uuid при форматировании хранилищ всех брокеров. Повторная генерация отдельного идентификатора для каждого узла создаст разные кластеры, которые не смогут объединиться.

Запустите брокер:

bin/kafka-server-start.sh config/kraft/server.properties

Проверьте работоспособность, создав тестовый топик:

bin/kafka-topics.sh --create --topic test --partitions 1 --replication-factor 1 --bootstrap-server localhost:9092

Топик с replication-factor 1 подходит для локальной проверки, но не для production: при отказе единственного брокера его партиция останется без копии.

Настройка кластера из нескольких брокеров

Для каждого брокера создайте отдельный файл server.properties. Ключевые параметры для трёх узлов:

  • broker.id - уникальный идентификатор в классической конфигурации; в KRaft идентичность узла задаётся параметром node.id
  • listeners - адрес и порт прослушивания, например PLAINTEXT://:9092 для первого брокера, PLAINTEXT://:9093 для второго
  • advertised.listeners - адреса, которые Kafka передаёт клиентам и другим брокерам; в production здесь должны быть достижимые имена или IP, а не localhost
  • log.dirs - путь к каталогу хранения данных
  • num.partitions - количество партиций по умолчанию для новых топиков
  • default.replication.factor - фактор репликации по умолчанию

В режиме KRaft все брокеры используют общий идентификатор кластера. Запустите каждый брокер со своим конфигурационным файлом. Проверьте, что брокеры видят друг друга, через команду kafka-broker-api-versions.sh --bootstrap-server localhost:9092. В реальной сети указывайте адрес брокера, доступный с узла, где выполняется CLI.

Минимальная production-схема KRaft

Для небольшого кластера с совмещёнными ролями пример конфигурации первого узла может выглядеть так:

node.id=1
process.roles=broker,controller
listeners=PLAINTEXT://broker-1:9092,CONTROLLER://broker-1:9093
advertised.listeners=PLAINTEXT://broker-1:9092
controller.listener.names=CONTROLLER
listener.security.protocol.map=CONTROLLER:PLAINTEXT,PLAINTEXT:PLAINTEXT
controller.quorum.voters=1@broker-1:9093,2@broker-2:9093,3@broker-3:9093
inter.broker.listener.name=PLAINTEXT
log.dirs=/var/lib/kafka-1

На broker-2 и broker-3 должны отличаться node.id, hostname в listeners и advertised.listeners, а также путь log.dirs. Например, для второго узла используются node.id=2, broker-2 и /var/lib/kafka-2. Параметр controller.quorum.voters, идентификатор кластера, имена listener-ов и межброкерный протокол должны совпадать на всех узлах.

Для раздельных ролей на controller-узлах задают process.roles=controller, а на broker-узлах - process.roles=broker. В этом случае controller-узлы не обслуживают пользовательские партиции, а broker-узлы не участвуют в хранении журнала метаданных как контроллеры. Перед запуском проверьте DNS, доступность controller-портов и отсутствие одинаковых node.id.

Типовые настройки production-кластера

Универсального server.properties для любого кластера нет: значения зависят от размера сообщений, throughput, числа партиций, дисков и требований к задержке. Для кластера с replication.factor=3 базовый профиль может начинаться со следующих настроек:

min.insync.replicas=2
unclean.leader.election.enable=false
num.network.threads=8
num.io.threads=16
log.segment.bytes=1073741824
log.retention.ms=604800000
log.retention.bytes=-1
message.max.bytes=1048588
  • min.insync.replicas - минимальное число реплик в ISR, необходимое для записи при acks=all. При RF=3 значение 2 позволяет пережить отказ одной реплики, но блокирует запись, если в ISR останется меньше двух копий.
  • unclean.leader.election.enable=false запрещает выбирать лидером реплику, которая не входила в ISR. Это снижает риск потери подтверждённых данных, но при полном отсутствии синхронизированных реплик повышает вероятность недоступности партиции.
  • acks=all задаётся на стороне producer. Для критичных данных его указывают явно и проверяют вместе с min.insync.replicas, чтобы producer не считал запись успешной при недостаточном количестве синхронизированных копий.
  • retention.ms и retention.bytes задают срок и объём хранения на уровне топика. Обычно срабатывает ограничение, которое достигается раньше. На уровне брокера используются значения по умолчанию log.retention.ms и log.retention.bytes, а настройки конкретного топика имеют приоритет.
  • message.max.bytes ограничивает размер сообщения или batch на брокере. Если разрешены крупные записи, согласуйте значение с max.request.size producer и fetch.max.bytes и max.partition.fetch.bytes consumer.
  • num.network.threads и num.io.threads определяют рабочие потоки обработки сетевых и дисковых запросов. Увеличение значений без проверки CPU и времени ожидания диска может ухудшить ситуацию, поэтому параметры меняют после нагрузочного тестирования.
  • log.segment.bytes влияет на частоту ротации сегментов. Меньшие сегменты ускоряют применение retention, но увеличивают количество файлов и операций с индексами; большие сегменты уменьшают число файлов, но могут увеличивать время очистки и восстановления.

Пример настройки retention для отдельного топика:

bin/kafka-configs.sh --bootstrap-server localhost:9092 --alter --entity-type topics --entity-name my-topic --add-config retention.ms=604800000,retention.bytes=107374182400

Числа в примере являются стартовыми, а не универсальной рекомендацией. После изменения параметров проверьте latency producer и consumer, загрузку CPU, дисковый iowait, размер очередей запросов и состояние ISR.

Работа с CLI-инструментами Kafka

CLI-утилиты входят в дистрибутив Kafka и позволяют управлять кластером без написания кода. Основные инструменты: kafka-topics.sh, kafka-console-producer.sh, kafka-console-consumer.sh.

Создание и управление топиками с помощью kafka-topics.sh

Создание топика с тремя партициями и фактором репликации 2:

bin/kafka-topics.sh --create --topic my-topic --partitions 3 --replication-factor 2 --bootstrap-server localhost:9092

Просмотр списка топиков:

bin/kafka-topics.sh --list --bootstrap-server localhost:9092

Детальное описание топика с лидерами партиций и ISR:

bin/kafka-topics.sh --describe --topic my-topic --bootstrap-server localhost:9092

Удаление топика:

bin/kafka-topics.sh --delete --topic my-topic --bootstrap-server localhost:9092

Отправка и чтение сообщений: kafka-console-producer.sh и kafka-console-consumer.sh

Запуск продюсера для отправки сообщений в топик:

bin/kafka-console-producer.sh --topic my-topic --bootstrap-server localhost:9092

После запуска вводите сообщения с клавиатуры, каждое на новой строке. Для выхода нажмите Ctrl+C.

Запуск консьюмера для чтения всех сообщений с начала топика:

bin/kafka-console-consumer.sh --topic my-topic --from-beginning --bootstrap-server localhost:9092

Для указания группы потребителей добавьте параметр --group:

bin/kafka-console-consumer.sh --topic my-topic --group my-group --from-beginning --bootstrap-server localhost:9092

Проверить состояние групп потребителей и их отставание можно командой kafka-consumer-groups.sh --describe --group my-group --bootstrap-server localhost:9092.

Как проверить кластер Kafka

Проверку начинайте с метаданных топиков, а затем переходите к состоянию consumer group и controller quorum. Для базовой диагностики используйте следующие команды:

bin/kafka-topics.sh --describe --bootstrap-server localhost:9092
bin/kafka-topics.sh --describe --under-replicated-partitions --bootstrap-server localhost:9092
bin/kafka-topics.sh --describe --unavailable-partitions --bootstrap-server localhost:9092
bin/kafka-consumer-groups.sh --bootstrap-server localhost:9092 --describe --all-groups
bin/kafka-broker-api-versions.sh --bootstrap-server localhost:9092

Команда kafka-topics.sh --describe показывает для каждой партиции лидера, список всех реплик и ISR. Значение Leader: -1 означает, что лидер не назначен. Если список Isr короче списка Replicas, часть реплик отстаёт. Пустой результат команды с --under-replicated-partitions означает, что на момент проверки under-replicated partitions не обнаружены.

Команда kafka-consumer-groups.sh --describe --all-groups позволяет найти группы с растущим lag, проверить текущий offset и сравнить его с log end offset. Для отдельной consumer group используйте параметр --group my-group. В production запускайте команды с реальным адресом bootstrap-сервера и, если требуется, с параметром --command-config для аутентификации.

В KRaft состояние controller quorum проверяют отдельной командой:

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

В выводе проверьте наличие текущего лидера метаданных, список voters и отсутствие отставания followers. Команда относится к KRaft и не заменяет проверки брокеров, топиков и consumer group.

Как понять, что Kafka здорова

Kafka-кластер можно считать здоровым, если одновременно выполняются несколько условий:

  • все ожидаемые брокеры отвечают на запросы, а client и inter-broker listeners доступны по корректным адресам;
  • у каждой партиции есть лидер, а значение Leader не равно -1;
  • UnderReplicatedPartitions и OfflinePartitions равны нулю, а состав ISR соответствует требованиям replication factor;
  • controller quorum в KRaft имеет активного лидера и синхронизированные voters;
  • lag каждой consumer group не растёт дольше установленного SLA, а ребалансировки не происходят постоянно;
  • на дисках есть запас свободного места, а CPU, RAM, сеть и дисковый iowait не находятся в устойчивой зоне насыщения.

Отказоустойчивость и репликация в Kafka

Репликация - основной механизм защиты данных в Kafka. Каждая партиция имеет несколько копий, размещённых на разных брокерах. Количество копий задаётся параметром replication.factor.

Репликация партиций: лидер и последователи

Одна из реплик партиции назначается лидером. Лидер принимает все записи от производителей и отдаёт данные потребителям. Остальные реплики - последователи - реплицируют данные лидера. Они не обрабатывают клиентские запросы, только синхронизируются.

Список синхронизированных реплик называется ISR (in-sync replicas). Если последователь отстаёт от лидера больше допустимого порога, он исключается из ISR. При отказе лидера новый лидер выбирается только из реплик, входящих в ISR. Это гарантирует, что новый лидер имеет все подтверждённые записи.

Обработка сбоев: как Kafka сохраняет данные

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

Данные не теряются, если фактор репликации больше единицы и производитель использует подтверждение acks=all. При этом запись считается успешной только после того, как необходимое количество синхронизированных реплик сохранило сообщение. На практике эту гарантию усиливают параметром min.insync.replicas. Например, при RF=3 и min.insync.replicas=2 отказ одной реплики не останавливает запись, а потеря двух реплик приводит к отказу новых записей вместо подтверждения недостаточно надёжной операции.

Параметр unclean.leader.election.enable=false не позволяет назначить лидером реплику вне ISR. Это сохраняет согласованность подтверждённых данных, но при отсутствии доступной ISR-партиции может привести к её временной недоступности. Если фактор репликации равен 1, отказ брокера приводит к потере данных.

Практические советы по эксплуатации Kafka

Эксплуатация production-кластера требует мониторинга и готовности к типичным проблемам. Ниже - ключевые метрики, стартовые пороги и частые неисправности.

Мониторинг кластера Kafka

Kafka отдаёт метрики через JMX. Ключевые показатели для отслеживания:

  • UnderReplicatedPartitions - количество партиций с неполной репликацией. Значение больше нуля не должно сохраняться дольше короткого окна после планового перезапуска; устойчивое значение больше нуля в течение 1-5 минут требует проверки.
  • OfflinePartitions - количество партиций без лидера. Должно быть равно нулю. Любое ненулевое значение - повод для немедленного разбора доступности партиций.
  • ConsumerLag - отставание потребителей. Нулевой lag не обязателен, но его постоянный рост в течение 5-10 минут или превышение допустимой задержки по SLA указывает, что consumer group не успевает обрабатывать поток.

Дополнительно отслеживайте состояние инфраструктуры:

  • свободное место на дисках брокеров: запас менее 15-20% можно использовать как предупредительный порог, менее 10% - как критический, с поправкой на скорость роста данных и размер сегментов;
  • дисковый iowait, latency и время обработки запросов: их рост вместе с UnderReplicatedPartitions часто указывает на проблему с диском или перегруженный storage;
  • сетевой throughput, ошибки интерфейсов, время ожидания request queue и NetworkProcessorAvgIdlePercent: устойчивое снижение idle ниже 20% может означать насыщение сетевых потоков;
  • CPU и паузы Java GC: устойчивые значения CPU выше 80-90% или длинные паузы могут увеличивать задержку репликации и приводить к росту lag;
  • RAM и page cache: низкий объём свободной памяти сам по себе не означает проблему, потому что Kafka использует RAM для page cache, но активный swap и OOM требуют немедленной проверки.

Для сбора метрик используйте Prometheus с JMX-экспортером, для визуализации - Grafana. Пороговые значения нужно привязать к базовой линии конкретного кластера и бизнес-SLA, а не задавать только по абсолютным числам. Подробная настройка мониторинга и алертинга описана в руководстве по нагрузочному тестированию и мониторингу Kafka.

Симптомы, причины и первичная диагностика

СимптомВероятная причинаЧто проверить
Нет лидера у партицииБрокер недоступен, ISR пуста, нарушена связь с controller quorum или неверно настроены listenerskafka-topics.sh --describe --topic my-topic, состояние брокеров, логи в logs/, доступность inter-broker и controller-портов
Растёт lagМедленная обработка сообщений, слишком мало партиций, частые ребалансировки или блокирующие операции в consumerkafka-consumer-groups.sh --describe --group my-group, processing latency, records-lag-max, состав consumer group и загрузку базы данных
Есть under-replicated partitionsОстановлен брокер, не хватает дисковой производительности, заполнен диск или есть потери пакетовkafka-topics.sh --describe --under-replicated-partitions, свободное место, iowait, сетевые ошибки, логи follower fetcher
Не создаётся топикНедоступен controller quorum, включены ограничения ACL, replication factor больше числа брокеров или некорректны listenersтекст ошибки CLI, логи брокера, kafka-metadata-quorum.sh в KRaft, число доступных брокеров и значения default.replication.factor

Типичные проблемы и их решение

Брокер не запускается. Частая причина - нехватка дискового пространства, неправильный broker.id или node.id, конфликт порта и ошибка в конфигурации KRaft. Проверьте логи брокера в каталоге logs/, убедитесь, что log.dirs указывает на доступный раздел с достаточным свободным местом, а идентификатор узла уникален.

Топик недоступен. Проверьте состояние партиций командой kafka-topics.sh --describe --topic my-topic --bootstrap-server localhost:9092. Если лидер партиции отсутствует, убедитесь, что все нужные брокеры запущены, ISR не пуста, controller quorum доступен и сеть между узлами работает.

Потребитель отстаёт. Увеличьте количество партиций или добавьте потребителей в группу только после проверки причины. Если все потребители заняты, а партиций недостаточно, масштабирование группы не поможет. Также проверьте, что обработка сообщений не блокируется медленными операциями, например синхронными запросами к базе данных.

Для развёртывания Kafka-кластера на облачной инфраструктуре подойдут облачные серверы Timeweb Cloud с настраиваемыми ресурсами. Если вы планируете строить event-driven систему на Kafka, изучите практическое руководство по настройке Kafka как Event Bus.

Типовые ошибки новичков

  • Выбирать слишком мало партиций, рассчитывая позднее просто добавить consumer-ов. Количество активных потребителей в одной consumer group ограничено числом партиций, поэтому будущий рост нужно учитывать заранее.
  • Использовать replication.factor=1 в production или задавать replication factor больше числа доступных брокеров. Первый вариант не даёт отказоустойчивости, второй приводит к ошибке создания топика.
  • Не задавать acks=all для критичных событий и не согласовывать его с min.insync.replicas. Проверяйте фактическую конфигурацию producer, особенно после миграций и смены библиотек.
  • Пытаться уменьшить число партиций. Kafka не поддерживает такое изменение напрямую: обычно создают новый топик с нужной схемой и выполняют контролируемую миграцию данных.
  • Использовать localhost в advertised.listeners на многосерверном кластере. Клиенты и другие брокеры будут получать недоступный адрес, даже если сам порт локально открыт.
Поделиться:
Сохранить гайд? В закладки браузера