System design · Брокеры сообщений
Как работает Apache Kafka: топики, партиции, брокеры и KRaft
Архитектура Kafka по частям: лог, партиции, репликация, контроллеры, consumer groups, транзакции, tiered storage, Kafka Streams, Kafka Connect и Schema Registry. Для каждой части: что это, как работает и зачем нужна.
Коротко
Apache Kafka — распределённая pub/sub-система обмена сообщениями с открытым исходным кодом. В её основе лежит лог: журнал, в который записи только дописываются в конец и читаются по порядку. Продюсеры пишут данные в Kafka один раз, а любое число консьюмеров читает их независимо и в реальном времени.
- топик делится на партиции, каждая партиция — отдельный лог, поэтому кластер масштабируется добавлением брокеров;
- каждая партиция обычно хранится в трёх копиях на разных брокерах, а пишет в неё только лидер, поэтому отказ диска или узла не приводит к потере данных;
- метаданными кластера управляет активный контроллер, которого выбирает алгоритм KRaft; начиная с Kafka 4.0 ZooKeeper не нужен;
- сообщения хранятся заданный срок (обычно 7 дней), а не до прочтения, поэтому историю можно перечитать;
- consumer groups делят партиции между консьюмерами, а Kafka Streams, Kafka Connect и Schema Registry добавляют потоковую обработку, интеграции и схемы.
Что такое Apache Kafka?
Apache Kafka — распределённая система обмена сообщениями по модели publish/subscribe с открытым исходным кодом. Она надёжно хранит потоки сообщений на диске, масштабируется горизонтально добавлением узлов, переживает отказы серверов и поставляется с инструментами для интеграций и потоковой обработки.
- Создана: LinkedIn, 2010
- Open source: с 2011, Apache Software Foundation
- Протокол: собственный, поверх TCP
- Основа: лог
В любой статье о Kafka её описывают одними и теми же словами:
«Это распределённая, надёжная, хорошо масштабируемая и отказоустойчивая pub/sub-система обмена сообщениями с открытым исходным кодом, богатыми интеграциями и возможностями потоковой обработки».
Всё верно, но человеку, который знакомится с Kafka впервые, такая фраза мало что объясняет. Ниже каждое слово определения разобрано по отдельности: что оно значит и за счёт какого механизма Kafka этого добивается. Сводная таблица — в разделе «Что означает каждое слово в определении Kafka?».
Кто создал Kafka и кто её развивает?
Kafka разработали в LinkedIn в 2010 году как систему обмена сообщениями. В 2011 году исходный код открыли и передали в Apache Software Foundation — отсюда open source в определении.
Сегодня это один из самых популярных open-source-проектов в инфраструктуре данных и стандартный инструмент дата-инженера. По открытым оценкам, Kafka пользуются более 70% компаний из списка Fortune 500 и 150 000 организаций.
Компанию, которая стоит за Kafka, знают под именем Confluent. Она стала публичной и расширила бизнес до полноценной платформы потоковой передачи данных поверх Kafka.
Kafka давно переросла простую систему обмена сообщениями: вокруг неё сложилась экосистема компонентов, которые вместе образуют платформу потоковой обработки данных (streaming platform). Её часто называют швейцарским ножом инфраструктуры данных.
Какую проблему решает Kafka?
Kafka решает задачу интеграции данных в масштабе: вместо отдельного конвейера между каждой парой сервисов все системы пишут данные в одно центральное хранилище и читают их оттуда через единый API. Число связей перестаёт расти квадратично вместе с числом сервисов.
LinkedIn нужно было связать между собой множество сервисов. Наивный способ — написать отдельную интеграцию «точка — точка» (такие интеграции называют конвейерами данных, data pipelines) между каждой парой сервисов. Это дало бы хаос порядка O(N²): конвейеры постоянно ломались бы, а при N в сотни их было бы невозможно поддерживать.

Kafka переворачивает задачу. Вместо конвейера на каждое соединение она предлагает:
- хранить данные в одном центральном месте — в Kafka;
- пользоваться одним стандартным API — Kafka API;
- подписывать приложения на эти данные и читать их в реальном времени.
Так пишущие системы отвязываются от читающих: первые просто публикуют данные в Kafka, вторые подписываются на них.
Как Kafka работает как pub/sub-система?
Pub/sub-модель в Kafka означает, что продюсеры публикуют сообщения в Kafka, а консьюмеры подписываются на них и читают независимо друг от друга. Одно и то же сообщение записывается один раз и читается столько раз, сколько нужно разным системам.
Данные надёжно сохраняются на диск на ограниченный срок, например на 7 дней. Kafka изначально проектировали под read fanout — сценарии, где одно сообщение читают несколько систем. Поэтому пропускная способность чтения в кластере обычно в несколько раз выше пропускной способности записи.

С Kafka не нужно поддерживать десятки хрупких самописных конвейеров, которые ломаются при перезапуске любой виртуальной машины. Данные записываются в Kafka один раз и читаются сколько угодно раз любой системой, которой они нужны.
Сценарии применения Kafka в этой статье дальше не разбираются. Зачем Kafka понадобилась LinkedIn, кратко сказано в разделе «Что такое Kafka Connect?».
Что такое лог в Kafka?
Лог — основная структура данных Kafka: упорядоченный журнал, в который записи можно только дописывать в конец. Удалять и изменять отдельные записи нельзя, а читают лог последовательно, слева направо, в том порядке, в котором записи добавлялись.
- Запись: только в конец (append-only)
- Стоимость записи: O(1)
- Адрес записи: offset
Kafka построена на простой структуре данных «лог». Всё остальное в системе — топики, репликация, метаданные, прогресс консьюмеров — устроено как логи.
Что такое offset в Kafka?
Offset — уникальный монотонно растущий номер записи в логе. Он указывает на конкретную запись и задаёт её место в порядке.

Какой API у лога?
API лога очень простой: дописать запись в конец и прочитать последовательный кусок между двумя offset.

Почему Kafka хранит лог на диске?
Kafka держит лог на диске, и последовательные операции лога хорошо ложатся на жёсткие диски (HDD). У HDD очень высокая пропускная способность при последовательном чтении и записи — в отличие от случайного доступа, с которым жёсткие диски справляются плохо.
Что такое сообщение в Kafka?
Сообщение в Kafka (запись, событие; record, message, event) — одна запись в логе. По сути это пара ключ-значение из сырых байтов: необязательный ключ byte[] key и значение byte[] value, плюс метаданные — offset, метка времени и пользовательские заголовки.
- Ключ: необязательный
- Формат: сырые байты
- Схема: на стороне клиента
Слова «запись», «сообщение» и «событие» в этой статье — синонимы. Ключ не обязателен: сообщение может состоять из одного значения.
![Класс Record на Java: сообщение Kafka состоит из необязательного ключа и значения, оба типа byte[]](../assets/blog/kafka/06.png?v=1e1b8779)
Главное, что нужно запомнить: ключ и значение — это сырые байты. Kafka сама по себе не поддерживает ни типы (int64, string и так далее), ни схемы — описания конкретной структуры сообщения.
Применять схемы — задача клиентского кода:
- при записи продюсер преобразует объекты в байты — сериализует их;
- при чтении консьюмер разбирает пришедшие по сети байты обратно в объект — десериализует их.
Что такое топик в Kafka?
Топик (topic) — именованная категория сообщений в Kafka, аналог таблицы в базе данных. Как для учётных записей пользователей и заказов в базе заводят разные таблицы, так в Kafka для разных потоков данных заводят разные топики.
- Аналог: таблица в базе данных
- Типично в кластере: от сотен до тысяч
- Состоит из: партиций
Одного лога мало: данные хочется разделять по категориям. Поэтому в кластере Kafka обычно от сотен до тысяч топиков.
Что такое партиция в Kafka?
Партиция (partition) — часть топика: Kafka шардирует каждый топик на одну или несколько партиций. Каждая партиция — отдельный экземпляр лога со своими offset, поэтому партиции одного топика можно хранить на разных брокерах и читать параллельно.
- Что это: шард топика
- Порядок: внутри партиции
- Типично: десятки на топик
Kafka — распределённая система, рассчитанная на масштаб, который не потянет одна машина. Поэтому она использует шардирование (sharding) — разбиение данных на части, которые хранятся на разных узлах.
Топик может состоять и из одной партиции, но обычно их десятки: так чтение распараллеливается (подробнее — в разделе о consumer groups).

Что такое продюсер и консьюмер в Kafka?
Продюсер и консьюмер (Producer и Consumer) — основные клиенты Kafka. Продюсер — класс, который записывает сообщения в топик. Консьюмер — класс, который подписывается на топик или конкретную партицию и читает из неё поток сообщений.
- Протокол: собственный, поверх TCP
- Официальная библиотека: Java
- Классы: KafkaProducer, KafkaConsumer
Kafka не использует HTTP — у неё собственный протокол поверх TCP. Значит, для отправки и получения запросов нужно больше специального кода: взять любую HTTP-библиотеку не получится.
Kafka поставляет собственные библиотеки, которые реализуют этот протокол. Проект Apache Kafka предлагает библиотеку на Java:
import org.apache.kafka.clients.producer.KafkaProducer;
import org.apache.kafka.clients.consumer.KafkaConsumer;Как продюсер отправляет сообщения в Kafka?
Класс Producer отправляет сообщения в топик. Партицию можно выбрать явно или доверить выбор Kafka.

Как консьюмер читает сообщения из Kafka?
Класс Consumer позволяет подписаться на топик (или конкретную партицию) и читать поступающий в него поток сообщений:

Это самая важная часть API, которую можно показать коротко; на деле методов гораздо больше. Kafka выглядит простой, но чтобы работать с ней эффективно, придётся разобраться во многих деталях.
Что такое брокер Kafka?
Брокер Kafka (broker) — экземпляр сервера Kafka, то есть узел распределённой системы. Все брокеры вместе образуют кластер. Kafka масштабируется горизонтально, добавлением узлов, поэтому обычное развёртывание состоит минимум из трёх брокеров.
- Брокер: один сервер Kafka
- Кластер: все брокеры системы
- Минимум: 3 узла
Именно это и значит распределённая в определении Kafka: система изначально рассчитана на работу на нескольких узлах, а не на одном сервере.
Как масштабируется Kafka?
Kafka масштабируется горизонтально: каждая партиция — независимый лог, партиции распределяются по брокерам, и пропускную способность поднимают добавлением брокеров. Запись в лог стоит O(1) и не требует блокировок, поэтому Kafka дописывает данные так быстро, как позволяет диск.
У Kafka много интересных оптимизаций производительности, но главная её сила — горизонтальная масштабируемость. Это и есть масштабируемая в определении.
Ключ к масштабированию — структура лога: записи просто дописываются в конец, их нельзя изменить или удалить по отдельности.
Сообщения внутри партиции независимы друг от друга: у них нет гарантий более высокого уровня вроде уникальных ключей. Поэтому блокировки почти не нужны, и Kafka дописывает данные в лог с той скоростью, которую выдерживает диск.
Каждая партиция — отдельный лог, а брокеры в кластер можно добавлять, поэтому масштаб ограничен только тем, сколько брокеров вы готовы добавить. Теоретически ничто не мешает кластеру, который принимает 50 ГиБ/с записи, вырасти вдвое, до 100 ГиБ/с.

Как работает репликация в Kafka?
Репликация в Kafka — хранение каждой партиции в нескольких копиях, репликах, на дисках разных брокеров. Число копий задаёт настройка replication factor; самое распространённое значение — три. Когда существует три копии данных, отказ одного диска не приводит к их потере.
- Настройка: replication factor
- Типично: 3 реплики
- Размещение: разные брокеры и зоны доступности
Партиция в Kafka живёт не на одном брокере — она реплицируется на несколько. При replication factor 3 у лога партиции три копии (реплики), и лежат они на дисках разных брокеров.
Реплицируют данные по многим причинам, и одна из главных — сохранность: при трёх копиях отказ одного диска не приводит к потере данных.
В современных облачных развёртываниях брокеры разносят по разным зонам доступности (availability zones). Это даёт очень высокую надёжность и доступность: даже если сгорит целый дата-центр, кластер Kafka продолжит работу. Это и есть надёжная (durable) в определении.
Что такое лидер партиции в Kafka?
Лидер партиции — реплика, которая в данный момент принимает все новые записи и служит источником истины для лога. Остальные реплики — фолловеры (followers): они постоянно копируют данные у лидера и готовы его заменить. Если брокер лидера отказывает, лидерство переходит к другому брокеру.
Как только в распределённой системе появляются копии данных, появляется и масса пограничных случаев. Синхронизировать новые данные сложно: копии должны совпадать, а система должна как-то договориться о последнем состоянии.
Такие задачи решает целый класс сложных алгоритмов — распределённый консенсус (distributed consensus).
Kafka использует простую модель репликации с одним лидером. В каждый момент одна реплика — лидер, две другие — фолловеры, то есть горячий резерв.
Новые записи принимает только лидер. Фолловеры активно реплицируют у него данные. Читать можно и с лидера, и с фолловеров (чтение с фолловеров включается в настройках).
Когда брокер выходит из строя, система это замечает, и другие брокеры берут на себя лидерство в партициях, которые вёл отказавший брокер. Так Kafka обеспечивает высокую доступность — это отказоустойчивая в определении.

Что такое __cluster_metadata в Kafka?
__cluster_metadata — служебный топик Kafka из одной партиции, в котором хранятся все изменения метаданных кластера. Каждая запись — одно событие кластера (дельта), и при воспроизведении всех записей по порядку любой узел детерминированно получает одно и то же текущее состояние кластера.
В распределённой системе все узлы должны сходиться во мнении о последнем состоянии кластера. Брокерам нужно согласовывать некоторые изменения метаданных, например выбор нового лидера. Это снова задача распределённого консенсуса; её полное решение слишком сложно для вводной статьи, поэтому здесь — только общая схема.
Kafka использует централизованную модель координации. И центральный координатор — не что иное, как… лог.
Все изменения метаданных Kafka надёжно записывает в специальный топик __cluster_metadata из одной партиции. Такое хранилище наследует все преимущества топиков: отказоустойчивость, надёжность и, что для метаданных важнее всего, порядок.

Иначе говоря, партиция топика метаданных — источник истины о последних метаданных в Kafka.
На этот топик подписан каждый брокер кластера. В реальном времени брокер забирает последние закоммиченные обновления и применяет каждую новую запись к метаданным в памяти. Так складывается его представление о текущем состоянии кластера.
Если каждый брокер — фолловер этой партиции, возникает естественный вопрос: кто лидер? Какой узел решает, какие метаданные записать?
Что такое контроллер в Kafka?
Контроллер Kafka — особый брокер, который не хранит обычные топики, а управляет кластером, то есть служит его control plane. Обычно разворачивают три контроллера; в каждый момент активен один — лидер лога __cluster_metadata, остальные работают горячим резервом.
- Роль: control plane кластера
- Обычно: 3 контроллера
- Активный: один, лидер лога метаданных

Писать в лог метаданных может только активный контроллер. Он принимает все решения о метаданных кластера: выбирает новых лидеров партиций при отказе брокера, создаёт топики, меняет конфигурацию на лету и так далее.
Самое важное — активный контроллер определяет, какие брокеры живы. Каждый брокер отправляет ему heartbeat — периодический сигнал «я жив». Если брокер не присылает heartbeat 6 секунд подряд, контроллер отключает его от кластера (fencing) и назначает другие брокеры лидерами партиций, которые вёл отключённый брокер.

Внимательный читатель спросит: если активный контроллер выбирает лидеров партиций, то кто выбирает лидера самой __cluster_metadata? Эта партиция особая: её лидера выбирает собственный алгоритм распределённого консенсуса — KRaft.
Что такое KRaft в Kafka?
KRaft (Kafka Raft) — собственный алгоритм консенсуса Kafka, основанный на Raft. Контроллеры образуют кворум KRaft, который выбирает активного контроллера и подтверждает изменения лога метаданных: запись считается закоммиченной, только когда её сохранило большинство кворума.
- Основа: Raft
- Участники: кворум контроллеров
- Заменил: ZooKeeper
Выбор лидера в распределённой системе — частный случай задачи консенсуса. Алгоритмов консенсуса много: Raft, Paxos, Zab и другие. Kafka использует собственный алгоритм KRaft, вдохновлённый Raft.
У KRaft две ключевые задачи:
- Выбрать активного контроллера.
- Узлы-контроллеры образуют кворум Raft.
- Кворум по протоколу выборов Raft выбирает лидера партиции
__cluster_metadata. Лидер этой партиции и есть активный контроллер.
- Договориться о последнем состоянии лога метаданных.
- Обновления метаданных сначала дописываются в Raft-лог на активном контроллере.
- Закоммиченными они считаются, только когда их сохранило большинство кворума.
Лидеров всех остальных, обычных партиций определяет активный контроллер. Он записывает решение в лог метаданных, и после коммита кворумом контроллеров решение окончательно.
Итого выбор лидера в Kafka устроен в два уровня:
- лидера среди контроллеров (активного контроллера) выбирает вариант Raft — KRaft;
- лидеров среди обычных брокеров выбирает контроллер.
Чем KRaft отличается от ZooKeeper?
KRaft появился в Kafka относительно недавно. Много лет Kafka работала с ZooKeeper. Тогда контроллер был один: он выполнял те же задачи, что и сегодня, но, что важно, одновременно был и обычным брокером. Свои решения он сохранял в ZooKeeper, который внутри использовал алгоритм консенсуса Zab. Начиная с Kafka 4.0 режим ZooKeeper удалён, и KRaft — единственный вариант.
Как выбор лидера устроен в других системах?
Модель выбора лидера через координатора отличает Kafka от других систем. Например, Redpanda держит отдельный кворум Raft на каждую партицию.

Как долго Kafka хранит сообщения?
Kafka хранит сообщения заданный срок, обычно 7 дней, независимо от того, прочитал их кто-то или нет. По истечении срока хранения (retention) сообщение удаляется автоматически. Поэтому старые сообщения можно перечитать и обработать заново — это называют replayability.
Одна из ключевых идей в дизайне Kafka — дать возможность воспроизводить исторические данные и отвязать срок хранения от клиентов. Альтернативные системы обмена сообщениями хранят сообщение, пока его никто не прочитал, а после чтения удаляют.
Kafka переворачивает эту модель: политика хранения — простое правило по времени. Сообщение удаляется автоматически, если пролежало на брокере дольше заданного срока, обычно 7 дней. Это возможно потому, что производительность лога O(1) не падает с ростом объёма данных.
Replayability — возможность заново обработать старые сообщения — крайне полезна. Например, в консьюмере долго жил незамеченный баг, и он неправильно обрабатывал сообщения. Когда баг исправлен, правильную логику можно прогнать по тем же сообщениям ещё раз.
Что такое tiered storage в Kafka?
Tiered storage (многоуровневое хранение) — механизм Kafka, который выгружает старые данные с дисков брокеров во внешнее объектное хранилище, например в S3. На брокерах остаются только свежие «горячие» данные, а историю по-прежнему читают через обычный API Kafka.
- Горячий уровень: диски брокеров
- Холодный уровень: S3 и другие облачные объектные хранилища
- Экономия: более чем в 10 раз (оценка)
В большом масштабе хранить столько исторических данных крайне сложно. Кластер с потоком записи 1 ГБ/с накапливает 1 772 ТБ данных (при хранении 7 дней и трёх репликах). Даже если разложить их на 100 брокеров, на диске каждого окажется 17 ТБ.
С таким объёмом состояния начинают копиться проблемы:
- Система теряет эластичность: любое действие или инцидент требует переноса огромного объёма данных, а это долго.
- Из-за устройства облачных цен самостоятельно хранить данные на HDD обычно выходит в 10 раз дороже, чем в S3.

Сообщество Kafka нашло остроумный способ решить все эти проблемы одной идеей — отдать их S3.
Это может звучать слишком просто или даже лениво, но решение очень элегантное. S3 — чудо инженерии: его поддерживают сотни сильных инженеров, и это, скорее всего, крупнейшая система хранения из известных человечеству.
Kafka сохраняет холодные данные на вторичный уровень хранения через подключаемый интерфейс. В качестве вторичного уровня поддерживаются объектные хранилища всех трёх крупных облаков, и интерфейс можно расширять.
Как устроен путь данных в Kafka с tiered storage?
- Горячий уровень (hotset tier). Сообщение записывается на брокер Kafka и реплицируется на реплики в кластере — оно хранится на дисках всех трёх узлов.
- Асинхронно сообщение выгружается в S3 — вторичный, холодный уровень.
- Холодный уровень (cold tier). Через настраиваемый срок (например, 12 часов) сообщение удаляется с брокеров, и единственным источником истины остаётся S3. Из S3 оно удаляется по истечении отдельного настраиваемого срока.

Холодные исторические данные по-прежнему читаются из Kafka обычным API. Разница только в том, что брокер достаёт их из S3, а не со своего диска.
Плюсы и минусы tiered storage
Минусы. Задержка при чтении исторических данных немного растёт, но это можно сгладить кэшированием.
Плюсы. Задержка на горячих данных может даже снизиться: становится экономически оправданно ставить быстрые SSD вместо HDD. Пропускная способность остаётся такой же высокой. Kafka становится гораздо эластичнее: при добавлении и удалении брокеров больше не нужно перемещать огромные объёмы данных. А хранить большие объёмы данных в Kafka становится более чем в 10 раз дешевле.

Что такое consumer group в Kafka?
Consumer group (группа консьюмеров) — набор экземпляров консьюмера, обычно на разных узлах, которые читают топик как одно целое и делят его партиции между собой. Каждая группа читает топик независимо и в своём темпе, поэтому один поток могут обрабатывать несколько разных приложений.
- Внутри группы: одна партиция — один консьюмер
- Между группами: чтение независимое
- Состав: меняется на лету
Напомним, что лог читается последовательно и по порядку. Отсюда набор требований:
- Топик делят на партиции, потому что данных много: один узел не должен справляться со всем топиком, нужно много консьюмеров.
- Партицию в каждый момент читает только один консьюмер — так порядок сообщений сохраняется без блокировок.
- Консьюмерам нужно договориться, как поделить партиции между собой.
- При этом Kafka должна позволять параллельно читать одни и те же партиции нескольким читателям — для сценариев с большим read fanout.
Kafka решает это через consumer groups: консьюмеры одной группы распределяют работу между собой и делят партиции, а разные группы читают топики независимо.
Состав группы динамический: чтение можно масштабировать вверх и вниз, добавляя и убирая участников прямо во время работы.
Пропускную способность чтения в Kafka можно наращивать двумя способами:
- Добавить консьюмеров в группу.
- Например, поток в топике вырос с 10 до 20 МБ/с, и два консьюмера не успевают. Добавьте ещё, чтобы они забрали дополнительную нагрузку.
- Добавить группы консьюмеров.
- Например, топик уже читают, но данные нужны ещё одному, отдельному приложению. Скажем, ночной бухгалтерской задаче, которая догоняет платежи за прошедший день.

Что такое Group Coordinator в Kafka?
Group Coordinator — брокер Kafka, который ведёт конкретную группу консьюмеров. Консьюмеры одной группы не общаются друг с другом напрямую: они шлют координатору heartbeat и по pull-модели узнают у него, какую работу выполнять. Решения принимает координатор — это централизованная модель координации.
Консьюмеры одной группы образуют распределённую систему обработки, а значит, снова упираются в задачи распределённых систем. Им нужно:
- договориться о прогрессе — до какого offset уже прочитано;
- следить за живостью — не отключился ли консьюмер и кто заберёт его работу;
- поддерживать динамический состав — не появился ли новый консьюмер;
- распределять работу — какой консьюмер какую партицию берёт.
Консьюмеры координируются косвенно, через брокер-координатор. Этот механизм называют протоколом членства в группе (consumer group membership protocol).
Что такое __consumer_offsets в Kafka?
__consumer_offsets — служебный топик Kafka, в котором группы консьюмеров хранят прогресс чтения: пары «партиция — offset», до которого сообщения уже обработаны. После сбоя консьюмер или его замена продолжает чтение с последнего сохранённого offset, а не с начала.
Group Coordinator выступает и «базой данных», которая хранит прогресс каждого консьюмера. Группы записывают простое соответствие {partition, offset} в специальный топик __consumer_offsets — так сохраняется, до какой записи они дочитали.
Прочитав сообщения, консьюмер коммитит через брокер-координатор offset, до которого обработал лог. Такие регулярные контрольные точки делают переключение при сбое плавным: консьюмер может перезапуститься и продолжить с того же места, или его работу подхватит другой консьюмер.
У топика __consumer_offsets много партиций, разбросанных по брокерам. Каждая группа привязана к определённой партиции, и лидер этой партиции работает координатором группы. Поэтому координатором какой-нибудь группы может быть любой брокер кластера — это исключает горячие точки, когда один брокер обслуживает все группы.

Протокол consumer groups — критически важная часть Kafka. Он достаточно общий, чтобы служить не только группам консьюмеров: он даёт внешним распределённым системам способ выбирать лидера и надёжно хранить состояние через Kafka.
На этом протоколе держатся и три системы, о которых речь пойдёт ниже: они используют его, чтобы работать как распределённые системы. Но сначала — о транзакциях.
Как работают транзакции и exactly once в Kafka?
Транзакция в Kafka позволяет продюсеру атомарно записать несколько сообщений в разные топики и партиции, даже на разных брокерах: консьюмеры увидят либо все эти сообщения, либо ни одного. В отличие от транзакций в базах данных, это прежде всего управление видимостью сообщений.
- Протокол: двухфазный коммит
- Координатор: Transaction Coordinator
- Exactly once: когда чтение и запись идут только через Kafka
Тема сложная, поэтому коротко. Транзакция в Kafka означает, что:
- продюсер может отправить много сообщений — в разные топики и партиции, на разные брокеры;
- эти сообщения атомарно либо коммитятся, либо отменяются на всех брокерах.
Технически это происходит с точки зрения консьюмера. Сообщения всё равно записываются в топики, но консьюмеры можно настроить так, чтобы они пропускали незакоммиченные (настройка isolation.level=read_committed). Это фильтрация на стороне клиента.
Коммит или отмена транзакции проходят по протоколу двухфазного коммита, и снова через централизованную модель координации: работу ведёт брокер — Transaction Coordinator. Устроено это довольно сложно.

Как Kafka защищает от дубликатов сообщений?
Главное в транзакциях — в типичных случаях они позволяют избавиться от дубликатов:
- Сбои сети или брокера. Если сеть потеряла ответ брокера или брокер перезапустился, тот же продюсер повторит запись идемпотентно — без дубликатов.
- Сбои самого продюсера. Если продюсер перезапускается с чистого состояния, он получает свой монотонно растущий идентификатор и увеличивает epoch — номер «поколения». Так старый экземпляр-«зомби» со старым epoch не сможет вмешаться в транзакцию.
Это не устраняет дубликаты полностью: пограничные случаи с внешними системами остаются. Но когда и чтение, и запись идут только через Kafka, без внешних систем, возможна обработка ровно один раз (exactly once processing). Этим активно пользуется Kafka Streams.
Из каких компонентов состоит Apache Kafka?
Проект Apache Kafka состоит из нескольких компонентов: Kafka Core — брокеры, контроллеры и координаторы; Kafka Clients — библиотеки продюсера и консьюмера; Kafka Streams — библиотека потоковой обработки; Kafka Connect — фреймворк интеграций с внешними системами. Schema Registry в сам проект не входит.
Первые два компонента разобраны выше, дальше — остальные. Все они живут в репозитории Apache Kafka на GitHub.
| Компонент | Что делает | Где работает | Часть Apache Kafka |
|---|---|---|---|
| Kafka Core | Хранит и реплицирует логи, управляет кластером, координирует группы и транзакции | Брокеры и контроллеры | Да |
| Kafka Clients | Записывает (Producer) и читает (Consumer) сообщения | В вашем приложении | Да |
| Kafka Streams | Непрерывно обрабатывает потоки: фильтры, join, агрегации по окнам | В вашем Java-приложении | Да |
| Kafka Connect | Переносит данные между Kafka и внешними системами через коннекторы | Отдельный кластер воркеров | Да |
| Schema Registry | Хранит схемы сообщений и отдаёт их клиентам | Отдельный HTTP-сервис | Нет, реализации сторонние |
Что такое Kafka Streams?
Kafka Streams — Java-библиотека для потоковой обработки данных на стороне клиента. Она непрерывно читает сообщения из входных топиков Kafka, обрабатывает их (map, filter, join, суммы, агрегации по окнам) и так же непрерывно пишет результат в выходной топик.
- Тип: клиентская библиотека, Java
- Вход и выход: только топики Kafka
- Гарантия: exactly once через транзакции
Потоковый обработчик (stream processor) — это приложение, которое:
- непрерывно обрабатывает сообщения;
- работает производительно и в масштабе;
- выдаёт результаты в реальном времени;
- поддерживает сложные операции с состоянием: агрегации по окнам и join нескольких потоков данных.
Kafka Streams — высокоуровневая библиотека клиентской потоковой обработки с богатым декларативным API. Простой пример на псевдокоде:

Этот код непрерывно считает сумму просмотров страниц людьми (без ботов) за последнюю минуту и пишет её в новый топик. Вот как при этом выглядят топики Kafka:

Как масштабируется Kafka Streams?
API предназначен для использования внутри ваших Java-приложений и работает как консьюмер. Kafka Streams распределяет работу по нескольким приложениям, как это делают consumer groups, но умеет распределять её ещё и по потокам (threads) внутри приложения. Для координации между экземплярами под капотом используется тот же протокол consumer groups.

Технически того же можно добиться своим кодом на простых библиотеках продюсера и консьюмера, но это очень много работы. Kafka Streams — абстракция более высокого уровня над обоими клиентами, с большим объёмом логики обработки, оркестрации и работы с состоянием сверху.
Как Kafka Streams обеспечивает exactly once?
Kafka Streams работает только с Kafka: берёт данные из одного топика и отправляет результат в другой. Поэтому она может гарантировать обработку ровно один раз через транзакции Kafka — на практике это значит, что данные обрабатываются атомарно.
Например, можно прочитать набор сообщений о платежах из топика, посчитать сумму и сохранить результат в другой топик со стопроцентной гарантией, что ни одно сообщение не потеряно и не посчитано дважды.
Если хотите подробнее — вот документация Kafka Streams. Именно Kafka Streams стоит за словами возможности потоковой обработки в определении.
Что такое Schema Registry в Kafka?
Schema Registry — внешний HTTP-сервис, который хранит схемы сообщений Kafka и отдаёт их клиентам. Сама Kafka типов и схем не знает, поэтому продюсеры регистрируют схему для топика, а консьюмеры скачивают её по ID из сообщения и с её помощью десериализуют данные.
- Протокол: HTTP
- Хранилище схем: топик Kafka
- Первая реализация: Confluent
Kafka не поддерживает ни типы (например, int64), ни схемы: сообщения — просто сырые байты. Но чтобы обработать сообщение, например сложить суммы оплат заказов, нужно разобрать его структуру и найти точное поле со значением. Как это сделать?
«Официального» способа нет: open-source-проект Apache схемы не поддерживает. Но есть распространённое соглашение — внешний HTTP-сервис с базой данных, в которой хранятся схемы.
На практике роль этой базы играет топик Kafka: он хранит простые пары {schema, topic}. Клиенты Kafka подключаются к сервису, скачивают схему и используют её при каждой сериализации и десериализации.
Первым таким реестром стал проект Schema Registry от Confluent с открытым исходным кодом, но не open-source-лицензией (source-available). Без официальной поддержки Apache и по-настоящему открытой лицензии экосистема раздробилась: сегодня есть много разных реализаций сервиса, а многие пользователи Kafka управляют схемами по-своему.
Как работает Schema Registry?
- Продюсеры выбирают схему, привязывают её к топику и регистрируют в реестре.
- Продюсеры сериализуют сообщение в нужной структуре, включая в него уникальный ID схемы, и пишут в Kafka.
- Консьюмеры получают сообщение, извлекают ID схемы и запрашивают (и кэшируют) схему из реестра.
- Консьюмеры десериализуют сообщение по этой схеме.

Что такое Kafka Connect?
Kafka Connect — фреймворк и среда выполнения для плагинов-коннекторов, которые переносят данные между Kafka и внешними системами. Для пользователя это no-code/low-code-способ подключить к Kafka, например, Elasticsearch, PostgreSQL или облачное аналитическое хранилище — в обе стороны.
- Source-коннектор: система → Kafka
- Sink-коннектор: Kafka → система
- Управление: REST API
Kafka создавали для решения задачи интеграции данных в LinkedIn — для переноса данных между разными системами. Это непросто: у каждой системы свой API, свой протокол (TCP, HTTP, JDBC и так далее), свой формат (XML, JSON, Protobuf, Avro) и свои гарантии совместимости.
Kafka Connect это стандартизирует. Connect — одновременно фреймворк (набор API) и среда выполнения для плагинов, которые связывают Kafka с внешними системами. Получается единый стандартный способ интеграции: сложный код, который гарантирует отказоустойчивость, порядок и обработку ровно один раз, пишется один раз в виде плагинов и проверен в бою.

Как устроен кластер Kafka Connect?
Три основных термина Connect:
- Connect Workers — обычные узлы, которые образуют распределённый кластер Connect.
- Connect Herder — воркер, который управляет кластером. Он предоставляет REST API, через который пользователи проверяют статус задач, запускают новые и так далее.
- Connectors — плагины (библиотеки), которые работают на воркерах и содержат код интеграции с другими системами.
- Source-коннектор читает данные из внешней системы и пишет их в Kafka (система → Kafka).
- Sink-коннектор читает данные из Kafka и пишет их во внешнюю систему (Kafka → система).
Пользователь поднимает кластер из нескольких воркеров, устанавливает на них jar-файлы нужного коннектора и запускает интеграцию простым HTTP-запросом POST.
Это ещё одна распределённая система обработки. Выбор лидера-Herder, общее членство в кластере и распределение новых задач по воркерам прозрачно выполняются через протокол consumer groups Kafka.
По сути Connect — это большой объём логики интеграции конкретных плагинов поверх обычных API KafkaProducer и KafkaConsumer. Существуют сотни плагинов-коннекторов, и именно они дают Kafka невероятно широкие возможности интеграции — это богатые интеграции в определении.
Что означает каждое слово в определении Kafka?
Определение Kafka — «распределённая, надёжная, хорошо масштабируемая и отказоустойчивая pub/sub-система обмена сообщениями с открытым исходным кодом, богатыми интеграциями и потоковой обработкой» — перечисляет восемь свойств. За каждым из них стоит конкретный механизм, разобранный в этой статье.
| Свойство | Что значит | За счёт чего |
|---|---|---|
| Open source | Исходный код открыт, проект развивает сообщество | Передача в Apache Software Foundation в 2011 году |
| Pub/sub-система | Пишут один раз, читают многие и независимо | Лог, топики, consumer groups |
| Распределённая | Работает на кластере из нескольких узлов | Брокеры, партиции |
| Масштабируемая | Пропускная способность растёт с числом брокеров | Партиции как независимые логи, запись O(1) без блокировок |
| Надёжная | Данные не теряются при отказе диска или дата-центра | Репликация, зоны доступности, tiered storage |
| Отказоустойчивая | Продолжает работать при отказе узлов | Лидеры и фолловеры, контроллер, heartbeat, KRaft |
| Потоковая обработка | Обработка данных в реальном времени | Kafka Streams, транзакции, exactly once |
| Богатые интеграции | Готовые подключения к внешним системам | Kafka Connect и сотни коннекторов |
Всё это Kafka получает благодаря множеству внутренних механизмов: библиотекам и API продюсера и консьюмера; топикам, партициям и репликам; брокерам и кворуму контроллеров KRaft; идемпотентности, транзакциям и exactly once; tiered storage; consumer groups и их протоколу; Kafka Streams, Kafka Connect и Schema Registry. Поэтому Kafka и называют швейцарским ножом дата-инжиниринга.
Как развивается Apache Kafka?
Apache Kafka — активный open-source-проект, который постоянно развивается. По состоянию на сентябрь 2025 года в работе были три заметные функции: очереди (Queues), diskless-топики, хранящие данные в объектном хранилище вместо дисков брокеров, и Iceberg-топики в открытом табличном формате.
- Очереди (Queues) — чтение партиции с семантикой очереди. Порядка у очередей нет, зато несколько консьюмеров могут читать один лог и подтверждать каждую запись по отдельности. Это отличается от классической модели Kafka «одна партиция — один консьюмер», где консьюмеры читают данные по порядку и знают только «я прочитал вплоть до этого сообщения».
- Diskless-топики — хранение партиций без лидера: слоем данных служит объектное хранилище (S3), а не диски брокеров. По оценкам, это может сократить облачные расходы на 90% и больше, ещё поднять масштабируемость и упростить эксплуатацию Kafka.
- Iceberg-топики — хранение данных Kafka в открытом табличном формате Iceberg без копирования (zero-copy).
Какие есть альтернативы Apache Kafka?
С Apache Kafka конкурирует множество проприетарных систем. Общее у них одно: все используют тот же протокол и API Kafka, просто реализуют их по-другому. Среди них — управляемые Kafka-сервисы облачных провайдеров, Redpanda, WarpStream, BufStream, AutoMQ и другие. Область очень богатая и быстро развивается.
Когда использовать Kafka?
Kafka берут, когда одни и те же события нужно надёжно доставить нескольким системам, обработать поток данных в реальном времени или связать много сервисов без конвейеров «каждый с каждым». Для очереди задач с поштучным подтверждением сообщений классическая модель Kafka подходит хуже.
| Задача или симптом | Что использовать в Kafka |
|---|---|
| Много сервисов связаны конвейерами «точка — точка» | Центральный лог: продюсеры пишут в топики, получатели подписываются |
| Одни и те же события нужны нескольким системам | Отдельная consumer group на каждую систему |
| Консьюмеры не успевают за потоком | Больше консьюмеров в группе, но не больше, чем партиций в топике |
| Нужно заново обработать историю после исправления бага | Срок хранения (retention) и повторное чтение со старого offset |
| Хранить историю на дисках брокеров слишком дорого | Tiered storage с выгрузкой в S3 |
| Нужны агрегаты по потоку в реальном времени | Kafka Streams |
| Нужно перенести данные из базы в хранилище или поиск | Kafka Connect: source- и sink-коннекторы |
| Продюсеры и консьюмеры расходятся в формате сообщений | Schema Registry |
| Нельзя допустить дубликатов при обработке внутри Kafka | Идемпотентный продюсер и транзакции, exactly once |
Какие ошибки чаще всего допускают при работе с Kafka?
- Ждут порядка сообщений по всему топику, хотя порядок гарантирован только внутри партиции.
- Запускают в группе больше консьюмеров, чем партиций в топике: лишние консьюмеры простаивают.
- Считают Kafka очередью, которая удаляет сообщение после прочтения, — а оно хранится весь срок retention.
- Путают транзакции Kafka с транзакциями базы данных: в Kafka это управление видимостью сообщений.
- Считают exactly once сквозной гарантией, хотя с внешними системами дубликаты всё ещё возможны.
- Пишут в топик байты без договорённости о схеме, и консьюмеры ломаются при первом изменении формата.
Частые вопросы о Kafka
Kafka — это очередь сообщений или база данных?
Kafka — распределённый лог с моделью publish/subscribe. В отличие от классической очереди, сообщение не удаляется после прочтения: оно хранится заданный срок, обычно 7 дней, и его может прочитать любое число консьюмеров. Базой данных в привычном смысле Kafka тоже не является: запросов по полям нет, данные читаются последовательно по offset.
Чем Kafka отличается от RabbitMQ?
RabbitMQ — классический брокер очередей: сообщение доставляется консьюмеру и удаляется после подтверждения, а гибкую маршрутизацию дают exchange. Kafka хранит сообщения в логе заданный срок, их можно перечитать, а одни и те же данные независимо читают несколько consumer groups. Kafka берут для потоков событий и высокой пропускной способности, RabbitMQ — для очередей задач и сложной маршрутизации.
Нужен ли ZooKeeper для Kafka?
Начиная с версии Kafka 4.0 — нет. Метаданные кластера хранятся в служебном топике __cluster_metadata, а активного контроллера выбирает встроенный алгоритм консенсуса KRaft. До перехода на KRaft Kafka много лет хранила решения контроллера в ZooKeeper.
Гарантирует ли Kafka порядок сообщений?
Да, но только внутри одной партиции: в ней записи упорядочены по offset, и в каждой consumer group партицию читает один консьюмер. Порядка между партициями одного топика нет. Чтобы связанные сообщения шли по порядку, им дают одинаковый ключ: стандартный партиционер отправляет сообщения с одним ключом в одну партицию.
Сколько партиций делать в топике Kafka?
Число партиций ограничивает параллелизм чтения: в одной consumer group партицию читает один консьюмер, поэтому консьюмеров больше, чем партиций, ставить бессмысленно. Обычно топику дают десятки партиций с запасом на рост. Увеличить число партиций можно, уменьшить нельзя, а при увеличении меняется распределение ключей по партициям.
Есть ли в Kafka exactly once?
Есть, когда и чтение, и запись идут только через Kafka. Идемпотентный продюсер убирает дубликаты при сбоях сети и брокера, а транзакции атомарно записывают сообщения в несколько партиций. Kafka Streams использует это для обработки ровно один раз. С внешними системами дубликаты всё ещё возможны.
Сколько Kafka хранит сообщения?
Столько, сколько задано политикой хранения (retention), обычно 7 дней, независимо от того, прочитали сообщение или нет. С tiered storage старые данные выгружаются в объектное хранилище вроде S3 и хранятся там отдельный, более долгий срок, а читаются тем же API.
Можно ли работать с Kafka из C#, Go и Python?
Да. Официальная клиентская библиотека проекта Apache Kafka написана на Java, а для других языков есть сторонние клиенты, реализующие протокол Kafka: например, confluent-kafka-dotnet для C#, franz-go и confluent-kafka-go для Go, confluent-kafka-python для Python. Kafka Streams доступна только на JVM.
Источники
- Apache Kafka Documentation, Design — устройство лога, репликации и хранения.
- The Raft Consensus Algorithm.
- KIP-405: Kafka Tiered Storage.
- Aiven, 16 ways tiered storage makes Kafka better.
- KIP-932: Queues for Kafka и KIP-1150: Diskless Topics.
- Apache Kafka Documentation, Kafka Streams.
- Apache Kafka на GitHub.
Apache®, Apache Kafka®, Kafka и логотип Kafka — зарегистрированные товарные знаки или товарные знаки Apache Software Foundation в США и/или других странах. Их использование не означает одобрения со стороны The Apache Software Foundation.