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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

Партиция физически состоит из нескольких сегментов - файлов фиксированного размера. По умолчанию размер сегмента - 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

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

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

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

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

  • broker.id - уникальный идентификатор: 0, 1, 2
  • listeners - адрес и порт прослушивания, например PLAINTEXT://:9092 для первого брокера, PLAINTEXT://:9093 для второго
  • log.dirs - путь к каталогу хранения данных
  • num.partitions - количество партиций по умолчанию для новых топиков
  • default.replication.factor - фактор репликации по умолчанию

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

Работа с 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

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

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

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

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

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

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

Данные не теряются, если фактор репликации больше единицы и производитель использует подтверждение acks=all. При этом запись считается успешной только после того, как все синхронизированные реплики сохранили сообщение. Если фактор репликации равен 1, отказ брокера приводит к потере данных.

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

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

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

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

  • UnderReplicatedPartitions - количество партиций с неполной репликацией. Значение больше нуля указывает на проблемы с синхронизацией.
  • OfflinePartitions - количество партиций без лидера. Должно быть равно нулю.
  • ConsumerLag - отставание потребителей. Растущий лаг сигнализирует о том, что консьюмеры не успевают обрабатывать поток.

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

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

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

Топик недоступен. Проверьте состояние партиций командой kafka-topics.sh --describe. Если лидер партиции отсутствует, убедитесь, что все брокеры запущены и сеть между ними работает.

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

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

Поделиться:
Сохранить гайд? В закладки браузера