Kafka
Справочник: Справочник по Apache Kafka.
Kafka
По сути, для самой Kafka сообщение — это просто массив байтов (byte array). Она абсолютно не заботится о том, что внутри: текст, картинка или данные от датчика. Это ключевое решение разработчиков, которое делает Kafka такой гибкой.
Ваше сообщение для Kafka — это "конверт", в котором есть фиксированная служебная часть и само содержимое. Внутри служебной части содержится важная техническая информация:
- CRC32: Контрольная сумма для проверки целостности (не повредилось ли сообщение в пути).
- Magic Byte: Версия формата сообщения.
- Attributes: Флаги, например, тип сжатия (Gzip, Snappy, LZ4) или временная метка.
- Ключ и Значение: Это и есть те самые массивы байтов, где лежат ваши данные. Ключ может быть null.
Вручную разработчики эти байты не пишут - всю работу берут на себя клиентские библиотеки (например, официальная Java-библиотека или kafkajs для Node.js). Технически, разработчик дергает метод библиотеки, и она делает всю магию по упаковке и отправке.
Процесс выглядит так:
- Создание объекта: В коде вы создаете объект, например,
ProducerRecordв Java. Вы указываете топик, ключ и значение.
producer.send(new ProducerRecord<String, String>("my-topic", key, value));
- Сериализация: Самый важный шаг. В настройках продюсера вы указываете сериализатор (serializer)
StringSerializerпревратит ваш текст в байты.JsonSerializerпревратит ваш JSON-объект в байты.AvroSerializerилиProtobufSerializerпревратят структурированные данные в компактный бинарный формат. Вы просто указываете класс сериализатора в конфиге. Библиотека сама берет ваш объект и превращает его в байтовый массив согласно этому сериализатору.
props.put("value.serializer", "org.apache.kafka.common.serialization.StringSerializer");
Для передачи сообщений Kafka использует свой собственный бинарный протокол поверх TCP. Он не основан на HTTP или AMQP. Любой запрос в этом протоколе, в том числе и отправка сообщения (Produce Request), выглядит так:
- Размер сообщения (
MessageSize): 4 байта, которые говорят клиенту, сколько байт ему нужно прочитать дальше. - Тело запроса (
RequestMessage): Здесь уже лежит сам "конверт" с данными. Для отправки сообщения это будет "ProduceRequest", который содержит информацию о топике, партиции и само сообщение (набор байтов).
Когда вы вызываете producer.send(), библиотека делает следующее:
- Сериализует ваши данные.
- Упаковывает их в структуру, описанную протоколом.
- Отправляет этот бинарный пакет по TCP-соединению на нужный брокер.
Для данных используются форматы JSON, Apache Avro, Protocol Buffers.
Kafka может применяться в двух основных режимах:
- как очередь сообщений для асинхронной обработки (например, постобработка видео на YouTube)
- как система потоковой обработки для непрерывного получения данных в реальном времени.
В качестве очереди сообщений каждое сообщение обычно обрабатывается одним потребителем в группе, а в качестве потока одни и те же данные могут быть прочитаны несколькими независимыми группами потребителей или воспроизведены из любой точки журнала.
Углублённая статья для разработчиков, архитекторов и инженеров. Определения компонентов, первый запуск и сценарии применения — в базовой главе Apache Kafka — потоковая обработка данных. CLI и таблицы параметров — в Справочнике по Apache Kafka.
Play ITЗагрузка интерактивного демо…
Ниже — архитектура кластера, гарантии доставки, интеграция, эксплуатация и безопасность.
Архитектура Kafka
Базовые концепции распределённой очереди сообщений
Механизм обработки сообщений
Когда какое-то событие происходит (например, пользователь кликнул по рекламе или оплатил заказ), программа-отправитель (производитель) создает сообщение. Это как письмо, которое нужно отправить. В каждом таком "письме" обязательно есть содержимое — сами данные. А еще могут быть:
- Ключ — как адрес на конверте, который решает, в какой ящик попадет письмо;
- Временная метка — когда письмо отправлено;
- Заголовки — дополнительная информация, как обратный адрес.
Ключ сообщения играет определяющую роль в маршрутизации, поскольку алгоритм хеширования ключа определяет конкретный раздел для записи, при этом сообщения с одинаковым ключом гарантированно попадают в один раздел, сохраняя порядок обработки. Представьте, что у вас есть несколько почтовых ящиков (в Kafka это "разделы"). Алгоритм смотрит на ключ, вычисляет его "отпечаток" и решает, в какой ящик положить письмо.
Если вы не указали ключ, Kafka сама решает, куда положить сообщение. Это похоже на то, как курьер случайным образом разбрасывает письма по разным ящикам, чтобы ни один не переполнялся - то есть, для сообщений без ключа применяется стратегия распределения по умолчанию, которая в современных клиентах предусматривает временную отправку в один выбранный раздел для эффективного формирования пакетов, после чего клиент переключается на новый раздел.
Представьте, что вы читаете книгу с закладкой. В Kafka смещение (offset) — это и есть такая закладка. Каждое сообщение имеет свой номер позиции в ящике. Потребитель (программа, которая читает сообщения) регулярно сообщает Kafka: "Я дочитал до сообщения №42, дальше продолжу с №43". Это называется фиксацией смещения. Если потребитель сломался после обработки сообщения, но до того, как успел зафиксировать смещение, он перечитает это сообщение заново после перезапуска. Такой подход гарантирует, что сообщение не будет потеряно, но иногда может быть обработано дважды. Это называют доставкой "хотя бы один раз" (at-least-once).
Смещения представляют собой уникальные идентификаторы, определяющие позицию сообщения в разделе, и используются потребителями для отслеживания прогресса чтения. Потребители периодически фиксируют свои смещения обратно в Kafka, что позволяет возобновить чтение с места остановки в случае сбоя или перезагрузки, причём стандартная схема с фиксацией после обработки обеспечивает доставку at-least-once.
Kafka предлагает простые способы повысить производительность:
- Отправлять сообщения пачками. Вместо того чтобы отправлять каждое письмо отдельно, вы собираете их в стопку и отправляете все сразу. Это как отнести сразу 10 конвертов на почту вместо 10 походов. Один вызов send или sendBatch может передать много сообщений.
- Сжимать сообщения. Перед отправкой данные можно упаковать архиватором (например, GZIP). Это как спрессовать одежду в вакуумный пакет — занимает меньше места и передается быстрее.
- Правильно выбирать ключ. Если ключ распределяет сообщения равномерно по разным ящикам, нагрузка ложится на все ящики одинаково. Это как если бы в очереди к 5 кассам в супермаркете люди распределялись равномерно — все работает быстро. Если ключ неудачный, один ящик переполняется, а другие пустуют — это создает "горячий раздел" и тормозит всю систему.
Kafka — это супер-эффективная почтовая система для программ. Она:
- Раскладывает письма по ящикам (разделам) на основе ключа;
- Гарантирует порядок для писем с одинаковым ключом;
- Помнит закладки (смещения), чтобы продолжить чтение после сбоев;
- Позволяет отправлять письма пачками и сжимать их для скорости;
- Требует правильного выбора ключа, чтобы все ящики загружались равномерно.
Топики и партиции
Топик (topic) — логический канал для хранения потока сообщений с общим назначением. Физически топик разделён на упорядоченные партиции (partitions), каждая из которых представляет собой неизменяемый упорядоченный журнал записей. Партиции обеспечивают горизонтальное масштабирование — одна партиция обрабатывается одним брокером, но топик может содержать множество партиций, распределённых по кластеру. Внутри партиции каждое сообщение имеет уникальный смещение (offset), определяющее его позицию в журнале.
Если проще, то топик — это как папка или канал, куда складываются сообщения одного типа. Каждая программа-отправитель кладет сообщения только в нужный топик. А программы-читатели подписываются только на те топики, которые им интересны. Например, система аналитики читает из топика "sales", а система оповещений — из топика "support".
Топик — это логическое понятие. Он просто объединяет сообщения по смыслу, но физически данные хранятся иначе. Всего лишь "виртуальная папка". Физически данные внутри топика разбиты на партиции.
Партиция — это настоящий ящик с сообщениями, который:
- Хранит сообщения строго по порядку (как очередь в магазине);
- Не позволяет изменять уже записанные сообщения (только добавлять новые в конец);
- Присваивает каждому сообщению уникальный номер — смещение (offset): 0, 1, 2, 3...
Можно провести аналогию:
- Топик = большая библиотека.
- Партиции = стеллажи в этой библиотеке.
Партиции ещё называют разделами - они являются фундаментальной единицей масштабирования и представляют собой упорядоченные неизменяемые последовательности сообщений, функционирующие по принципу файла журнала с операцией добавления. Именно благодаря разделам Kafka обеспечивает параллельную обработку сообщений, при этом порядок гарантируется исключительно внутри одного раздела, что критически важно для проектирования.
Топики называют темами - они служат логическим объединением разделов и позволяют организовывать публикацию и чтение данных. В отличие от разделов, которые являются физическими сущностями, темы представляют собой способ логической организации данных, при этом каждая тема может содержать один или несколько разделов, распределённых по разным брокерам.
Горячие разделы возникают при неравномерном распределении нагрузки, когда один раздел оказывается перегруженным. Для решения этой проблемы применяются такие стратегии, как случайное распределение без ключа (при допустимости потери порядка), добавление соли для улучшения распределения, использование составного ключа или применение обратного давления для замедления темпа поступления сообщений.
Каждая тема Kafka имеет настраиваемую политику хранения, определяющую срок хранения сообщений в журнале через параметры retention.ms и retention.bytes. По умолчанию сообщения хранятся 7 дней, однако этот период может быть увеличен для систем, требующих длительного хранения, с учётом возможного роста затрат на хранение и снижения производительности.
retention.ms— сколько миллисекунд хранить (например, 604800000 = 7 дней)retention.bytes— сколько байт хранить (например, 1 ГБ)
Когда срок истекает или размер превышен, старые сообщения автоматически удаляются. Это не база данных — Kafka не предназначена для вечного хранения. Если вам нужно хранить данные годами — можно увеличить срок, но тогда:
- Растут затраты на диски;
- Снижается производительность (больше данных нужно индексировать);
- Дольше восстановление после сбоев.
Продюсеры и потребители
Продюсер (producer) — клиентское приложение, публикующее сообщения в топик. При отправке продюсер определяет целевую партицию: явно по ключу сообщения (хеширование ключа гарантирует упорядоченность для одинаковых ключей) или через стратегию балансировки (например, round-robin). Потребитель (consumer) — приложение, считывающее сообщения из партиций. Потребитель отслеживает своё текущее смещение для каждой партиции, что позволяет возобновлять чтение с последней прочитанной позиции после перезапуска.
У продюсера есть два варианта:
- Указать ключ
{ key: "user_123", value: "Клиент совершил покупку" }
- Kafka хеширует ключ и определяет партицию;
- Все сообщения с одинаковым ключом всегда попадают в одну и ту же партицию;
- Это гарантирует порядок — если у вас 3 сообщения от user_123, они будут обработаны строго в той последовательности, в которой пришли.
- Не указывать ключ
{ value: "Какое-то событие" }
- Kafka сама решает, в какую партицию отправить;
- Обычно используется стратегия round-robin (по очереди: 1-я партиция, 2-я, 3-я, потом снова 1-я);
- Порядок сообщений не гарантируется.
Производители и потребители являются процессами, соответственно записывающими данные в темы и считывающими их оттуда, причём Kafka не анализирует содержимое сообщений, а только обеспечивает их сохранность и передачу.
На стороне производителя реализуются автоматические повторы с настраиваемым количеством попыток и интервалами между ними, при этом включение режима идемпотентного производителя предотвращает дублирование сообщений при повторных отправках. Представьте, что вы отправили письмо, но почтальон не смог его доставить (сервер временно недоступен). Вместо того чтобы сдаться, продюсер пытается снова:
{
retries: 5, // Попробует 5 раз
initialRetryTime: 100 // Подождёт 100 мс между попытками
}
Если после 5 попыток отправить не удалось — тогда уже выбрасывается ошибка. А идемпотентный продюсер гарантирует, что даже если он отправил одно и то же сообщение несколько раз (из-за сбоя сети или повторов), в Kafka оно попадёт только один раз. Представьте, что вы заказываете пиццу по телефону. Оператор не расслышал адрес, вы повторяете. Но если он запишет ваш адрес дважды, приедут две пиццы, а вы заплатите за две. Идемпотентность — это как если бы оператор понял, что заказ уже есть, и не дублировал его. Включается одной настройкой:
{ idempotent: true }
На стороне потребителя требуется разработка собственной логики обработки ошибок, включая создание отдельной темы для неудачных сообщений и выделенного потребителя для их обработки. Сообщения с многократными ошибками перемещаются в специальную очередь мёртвых сообщений для последующего анализа.
Потребитель — это программа, которая читает сообщения из Kafka. Он как получатель писем, который регулярно заглядывает в свой почтовый ящик. У потребителя есть смещение (offset) — это как закладка в книге. Он запоминает, какое сообщение прочитал последним, и при следующем запуске продолжает с того же места.
Consumer Group - это группа потребителей, которые работают вместе, читая один топик. Каждая партиция назначается только одному потребителю из группы. Представим, что у вас есть топик с 3 партициями и группа из 3 потребителей:
- Потребитель №1 читает партицию 0
- Потребитель №2 читает партицию 1
- Потребитель №3 читает партицию 2
Все работают параллельно, нагрузка распределена. Если потребителей меньше, чем партиций:
- Группа из 2 потребителей, а партиций 3
- Один потребитель читает 2 партиции, другой — 1
А если потребителей больше, чем партиций:
- Группа из 5 потребителей, а партиций 3
- 2 потребителя просто простаивают (им нечего читать)
Продюсер — отправитель. Кладёт сообщения в топик, выбирая партицию по ключу. Может переотправлять при ошибках и не создавать дубли.
Потребитель — получатель. Читает сообщения из партиций, запоминая позицию (смещение). Может работать в группе с другими потребителями.
Брокеры и кластеры
Брокер (broker) — отдельный узел Kafka, отвечающий за приём, хранение и доставку сообщений. Каждый брокер имеет уникальный идентификатор и управляет набором партиций (как лидеров, так и реплик). Кластер (cluster) — совокупность брокеров, работающих совместно под единым логическим пространством имён. Кластер обеспечивает отказоустойчивость через репликацию партиций и автоматическое перераспределение ролей при сбоях узлов.
Представьте, что это склад, на котором хранятся товары (сообщения). У каждого брокера есть свой номер (идентификатор), и он отвечает за:
- Приём сообщений от продюсеров;
- Хранение сообщений на диске;
- Отдачу сообщений потребителям.
Каждый брокер управляет несколькими партициями. Причём партиции могут быть двух типов:
- Лидер (Leader):
- Главная копия партиции;
- Именно через лидера происходят все записи и, по умолчанию, все чтения;
- Если лидер упал, его заменяет одна из реплик.
- Реплика (Replica):
- Резервная копия партиции;
- Постоянно копирует данные от лидера;
- Если лидер умирает, реплика становится новым лидером.
Один брокер может быть лидером для одних партиций и репликой для других — это распределяет нагрузку равномерно.
Брокеры представляют собой отдельные серверы (физические или виртуальные), образующие кластер Kafka, каждый из которых хранит данные и обслуживает клиентов. Увеличение количества брокеров позволяет масштабировать систему как по объёму хранимых данных, так и по числу обслуживаемых клиентов.
Кластер — это группа брокеров, которые работают вместе как единая система. Это как сеть складов одного логистического центра — все склады знают друг о друге и могут заменить друг друга в случае поломки.
- Если один брокер сломался, данные не потеряются — они скопированы на другие брокеры. Один из них станет новым лидером, и работа продолжится.
- Если данных становится слишком много для одного сервера, можно добавить новый брокер в кластер, и нагрузка распределится.
- Кластер продолжает работать, даже если несколько брокеров вышли из строя (в зависимости от настроек репликации).
В кластере всегда есть контроллер (Controller) — специальный брокер, который координирует всю работу:
- Следит, кто из брокеров жив, а кто умер;
- Решает, какая реплика станет новым лидером при сбое;
- Распределяет партиции между брокерами при добавлении новых узлов.
Ограничения отдельного брокера определяются такими факторами, как размер сообщений (рекомендуется менее 1 МБ для оптимальной производительности), характеристики дисков и сети, число разделов и настройки подтверждений. При необходимости масштабирования применяются два основных подхода: горизонтальное масштабирование с добавлением брокеров и стратегия разделения данных, причём ключевым фактором успеха является правильный выбор ключа сообщения для распределения по разделам.
Группы потребителей и балансировка нагрузки
Kafka использует модель вытягивания, позволяя потребителям самостоятельно запрашивать новые сообщения через определённые интервалы. Такой подход даёт потребителям контроль над скоростью потребления, упрощает управление ошибками, предотвращает чрезмерную загрузку медленных потребителей и обеспечивает возможность эффективного пакетного чтения.
Группа потребителей гарантирует, что каждое событие обрабатывается лишь одним потребителем внутри группы, а при выходе одного из потребителей из строя происходит автоматическое перераспределение разделов среди оставшихся участников группы.
Группа потребителей (consumer group) — логическая группа потребителей, совместно обрабатывающая сообщения из топика. Каждая партиция топика назначается ровно одному активному потребителю в группе. Это обеспечивает параллельную обработку: если топик содержит 12 партиций, максимум 12 потребителей в группе могут обрабатывать данные одновременно.
Балансировка нагрузки выполняется протоколом перебалансировки (rebalance protocol):
- Потребители регистрируются в группе через координатора группы (group coordinator) — специальный брокер, ответственный за управление состоянием группы.
- При изменении состава группы (добавление/удаление потребителя) запускается фаза перебалансировки.
- Выбирается лидер группы, который вычисляет план распределения партиций (partition assignment) с использованием стратегии (Range, RoundRobin, CooperativeSticky).
- План распространяется всем участникам; потребители останавливают обработку, применяют новое распределение и возобновляют чтение.
Стратегия CooperativeSticky (рекомендуемая в современных версиях) минимизирует перемещение партиций между потребителями при инкрементальных изменениях состава группы.
Static membership (group.instance.id) — consumer сохраняет "личность" при кратковременном рестарте. Кластер не запускает полный rebalance, если consumer вернулся в пределах session.timeout.ms. Это снижает "шторм" rebalance при деплоях.
Rebalance listener — callback ConsumerRebalanceListener вызывается до и после перераспределения партиций. Типичный паттерн: commit offset в onPartitionsRevoked, восстановление локального состояния в onPartitionsAssigned.
Внутреннее устройство кластера
Контроллер кластера
Контроллер (controller) — единственный брокер в кластере, ответственный за координацию метаданных:
- Назначение лидеров партиций при старте кластера или сбое узла
- Обработка запросов на создание/удаление топиков
- Управление ISR (In-Sync Replicas) — множеством реплик, синхронизированных с лидером
В режиме KRaft (Kafka Raft Metadata mode, с версии 3.3; единственный режим с Kafka 4.0) контроллер использует встроенный Raft-консенсус без ZooKeeper.
Механизм репликации
Каждая партиция имеет один лидер и набор реплик (обычно 2–3). Продюсеры пишут только в лидера; реплики асинхронно синхронизируются с лидером. Реплики, отставание которых не превышает порог replica.lag.time.max.ms, входят в множество ISR. Только реплики из ISR могут быть избраны новым лидером при сбое. Фактор репликации (replication factor) определяет общее количество копий данных для отказоустойчивости.
Репликация реализована по модели лидера и последователей, где каждый раздел имеет одного лидера, обрабатывающего все запросы на запись, и несколько реплик-последователей, расположенных на разных брокерах и копирующих данные с лидера. Контроллер кластера координирует весь процесс репликации, отслеживает состояние брокеров и управляет динамикой назначения лидеров и последователей, что обеспечивает непрерывную доступность раздела при отказе брокера.
Физическое хранение данных
Данные партиции хранятся на диске в виде сегментированных файлов журнала:
- Основной файл
.logсодержит последовательность сообщений фиксированного размера (по умолчанию 1 ГБ) - Файл индекса
.indexотображает смещение в позицию в файле для быстрого поиска - Файл временного индекса
.timeindexпозволяет искать сообщения по метке времени - Файл
.txnindexиспользуется для поддержки транзакционных операций
Сообщения могут подвергаться сжатию на уровне партиции:
none— без сжатияgzip,snappy,lz4,zstd— алгоритмы сжатия с разным балансом скорости/степениcompact— режим очистки (compaction), сохраняющий только последнее значение для каждого ключа (актуально для топиков-хранилищ состояний)
Сегменты неизменяемы после закрытия; новые сообщения пишутся в активный сегмент. Устаревшие сегменты удаляются по политике хранения (время или объём).
Гарантии надёжности Kafka
Kafka описывает поведение системы через явные гарантии — их можно сопоставить с ACID у реляционных СУБД, но семантика другая.
Гарантии платформы:
- Упорядоченность в партиции — если producer записал B после A в одну партицию, offset B больше offset A, и consumer прочитает B после A.
- Committed — producer считает запись успешной после записи во все in-sync реплики (при
acks=all) или по выбранному уровню ack. - Durability — committed-сообщение не теряется, пока жива хотя бы одна реплика партиции.
- Visibility — consumer по умолчанию читает только committed-записи; транзакционные "незавершённые" скрыты при
isolation.level=read_committed.
Настройки брокера для надёжности:
| Параметр | Рекомендация | Зачем |
|---|---|---|
replication.factor | ≥ 3 в prod | Отказ одного-двух брокеров |
min.insync.replicas | 2 при RF=3 | acks=all не подтвердит запись, если ISR < min |
unclean.leader.election.enable | false | Запрет выбора несинхронизированной реплики лидером (риск потери данных) |
log.flush.interval.messages | по умолчанию | Kafka полагается на OS page cache; явный fsync на каждое сообщение редко нужен |
Надёжный producer: acks=all, retries > 0, enable.idempotence=true (включает max.in.flight.requests.per.connection ≤ 5 и правильные retry).
Надёжный consumer: enable.auto.commit=false, commit offset после успешной обработки batch, идемпотентная бизнес-логика или upsert по ключу.
Семантика exactly-once
At-least-once допускает дубликаты при retry producer или повторном commit consumer. Для агрегаций и stream processing дубликат искажает результат — нужна семантика exactly-once. Общая модель слоёв и effectively exactly-once — Идемпотентность и семантика доставки.
Kafka реализует её двумя механизмами.
Идемпотентный producer (enable.idempotence=true):
- Broker присваивает producer ID (PID) и sequence number каждому сообщению в партиции.
- При retry broker отбрасывает дубликат с тем же sequence.
- Работает в пределах одной сессии producer; при
transactional.idPID сохраняется дольше. - Не защищает от дубликатов, если приложение отправляет "логически одно и то же" событие дважды с разными ключами.
Транзакционный producer (transactional.id):
- Группирует несколько записей (и commit consumer offset) в одну атомарную транзакцию.
- Consumer с
isolation.level=read_committedвидит сообщения только после commit транзакции. - Типичный сценарий: consume → transform → produce в другой топик + commit offset в одной транзакции (Kafka Streams делает это автоматически).
Транзакции Kafka гарантируют атомарность внутри кластера Kafka. Запись в PostgreSQL и Kafka в одной distributed transaction требует паттерна outbox или двухфазного подхода. Внешние side-effect (email, HTTP) всё равно нуждаются в идемпотентности на стороне получателя.
Kafka Connect
Kafka Connect — фреймворк для массовой интеграции без написания producer/consumer вручную. Коннекторы бывают source (внешняя система → Kafka) и sink (Kafka → внешняя система).
Когда Connect, когда свой клиент:
| Connect | Свой producer/consumer |
|---|---|
| Стандартный CDC (Debezium), репликация в DWH | Сложная бизнес-логика, фильтрация, enrichment |
| Много однотипных коннекторов (N таблиц → N топиков) | Низкая latency, тонкий контроль retry/DLQ |
| Единая операционная модель (REST API Connect) | Язык/стек без готового коннектора |
Пример топологии: MySQL (binlog) → Debezium source → топик orders → Elasticsearch sink для полнотекстового поиска.
Connect поддерживает Single Message Transforms (SMT) — лёгкие преобразования (маскирование поля, переименование) без отдельного stream job.
Принципы data pipeline (Kafka как буфер между системами):
- Своевременность — producer пишет в real-time, consumer может читать потоково или пакетами (раз в час).
- Надёжность — at-least-once на уровне Kafka; exactly-once сквозь pipeline — через idempotent sink, upsert или Connect offset API.
- Формат данных — договоритесь о схеме заранее (Avro + Schema Registry), иначе каждая пара producer/consumer "сшита" своим JSON.
Схемы событий и Schema Registry
Сырой JSON в топике быстро превращается в проблему — поля переименовываются, типы плывут, consumer'ы разных команд ломаются. Для контрактов событий обычно выбирают Avro, Protobuf или JSON Schema вместо самодельных binary-сериализаторов.
Почему Avro
- Схема описана отдельно (обычно JSON), payload — компактный binary.
- Обратная и совместимая эволюция — добавление поля с
default, удаление optional-поля; старые consumer'ы читают новые сообщения (новые поля → default/null). - Один формат для Java, Python, Go через codegen или
GenericRecord.
Пример эволюции: поле faxNumber заменили на email с "default": null — старые записи без email десериализуются с email=null, новые без faxNumber — с faxNumber=null.
Schema Registry
Реестр схем (часто Confluent Schema Registry, open source) хранит версии схем по subject (topic-value, topic-key). В Kafka уходит wire format: magic byte + schema id (4 bytes) + Avro payload — полная схема в каждое сообщение не кладётся.
[0x0][schema_id: 4 bytes][Avro binary data...]
Producer с KafkaAvroSerializer регистрирует схему при первой записи; consumer с KafkaAvroDeserializer подтягивает схему по id.
props.put("key.serializer", "io.confluent.kafka.serializers.KafkaAvroSerializer");
props.put("value.serializer", "io.confluent.kafka.serializers.KafkaAvroSerializer");
props.put("schema.registry.url", "http://schema-registry:8081");
Режимы совместимости
| Режим | Правило |
|---|---|
| BACKWARD (default) | Новая схема читает старые данные (добавление optional-полей) |
| FORWARD | Старые consumer'ы читают новые данные |
| FULL | Оба направления |
| NONE | Любые изменения (опасно в prod) |
В Java-приложение с Apache Kafka и PostgreSQL для обучения используется JSON + custom serializer. Для prod event contract — Avro/Protobuf + Schema Registry и проверка совместимости в CI.
Альтернативы без Confluent — Apicurio Registry, AWS Glue Schema Registry, protobuf + buf validate.
Kafka Streams — обзор
Kafka Streams — библиотека (JVM), а не отдельный кластер. Приложение читает топики, строит топологию (граф операций) и пишет результат обратно в Kafka.
Ключевые понятия из потоковой обработки:
- Таблично-потоковый дуализм — топик событий и compacted-топик "таблица состояния" взаимозаменяемы (changelog).
- State stores — локальные RocksDB для агрегаций; при rebalance состояние восстанавливается из changelog-топика.
- Окна — tumbling, hopping, session windows для агрегаций по времени.
- Гарантии — at-least-once по умолчанию; exactly-once через
processing.guarantee=exactly_once_v2.
Альтернативы — ksqlDB (SQL поверх Kafka), Flink, Spark Structured Streaming — выбор зависит от команды и требований к state и операционной модели.
Мониторинг и SLO
Consumer lag — главный индикатор "здоровья" pipeline: разница между последним offset в партиции и committed offset группы. Рост lag → consumer не успевает или упал.
Метрики брокера (JMX / Prometheus):
| Метрика | Что означает |
|---|---|
| Under-replicated partitions | Реплика отстаёт или брокер недоступен — риск потери данных при сбое |
| Offline partitions | Нет лидера — запись/чтение недоступны |
| Request handler idle % | Загрузка CPU брокера |
| Bytes in/out per sec | Пропускная способность |
Метрики producer: record-error-rate, request-latency-avg, размер batch.
Метрики consumer: records-lag-max, commit-latency-avg, частота rebalance.
SLO-пример: 99% событий обрабатываются с lag < 60 сек; алерт при URP > 0 дольше 5 мин.
Инструменты — AKHQ, Kafka UI, Confluent Control Center, Grafana + kafka_exporter, kafka-consumer-groups.sh --describe.
Выбор аппаратного обеспечения
Kafka дисковая система: последовательная запись в log-сегменты. Ориентиры по sizing:
- Диск — SSD/NVMe; отдельные тома под
log.dirs; RAID10 или JBOD (не RAID5 для write-heavy). - Память — page cache ОС критичен для read; 64 GB+ на брокер в высоконагруженных кластерах; heap JVM обычно 6–12 GB (больше — длинные GC-pause).
- Сеть — 10 GbE+ между брокерами и клиентами при больших объёмах репликации.
- CPU — умеренная нагрузка; сжатие (
lz4,zstd) и SSL увеличивают потребление.
В облаке — managed-сервисы (AWS MSK, Azure Event Hubs for Kafka, Confluent Cloud) снимают часть sizing, но партиционирование и RF остаются на стороне архитектора.
Зеркальное копирование между кластерами
Mirroring — репликация данных между отдельными Kafka-кластерами (в отличие от репликации партиций внутри одного кластера). Встроенный инструмент Apache Kafka — MirrorMaker 2 (MM2).
Когда нужно несколько кластеров:
| Сценарий | Суть |
|---|---|
| Региональные + центральный | Локальные кластеры в ЦОД/регионах; агрегированные топики зеркалируются в центральный для аналитики |
| HA / DR | "Горячий" или "холодный" standby-кластер с копией данных для аварийного переключения |
| Изоляция | Разные SLA, security boundary или команды — отдельные кластеры вместо одного "монстра" |
| Active-active | Два ЦОД принимают трафик; mirroring синхронизирует топики (сложнее: порядок, идемпотентность) |
Как работает MM2: consumer читает из source-кластера, producer пишет в target-кластер. Топики в target обычно получают префикс (например, source.orders). Offset'ы между кластерами не совпадают — это отдельные журналы.
Ограничения mirroring:
- Задержка репликации (секунды–минуты) — DR-кластер отстаёт от primary.
- Exactly-once не переносится "сквозь" mirroring автоматически.
- Active-active требует продуманного naming, compacted topics и идемпотентных consumer'ов.
Альтернативы MM2 — Confluent Cluster Linking, LinkedIn Brooklin, Uber uReplicator — для крупных multi-DC с тонкой настройкой lag и failover.
Безопасность Kafka
Безопасность Kafka строится на пяти столпах — аутентификация (кто вы), авторизация (что можно), шифрование (защита на проводе и at rest), аудит (кто что делал), квоты (лимит ресурсов). Система защищена настолько, насколько защищено самое слабое звено — включая клиентов, CI/CD и секреты.
Протоколы listener'ов
security.protocol | Шифрование | Аутентификация |
|---|---|---|
PLAINTEXT | Нет | Нет |
SSL | TLS | Опционально (mTLS) |
SASL_PLAINTEXT | Нет | SASL |
SASL_SSL | TLS + SASL | Рекомендуется для prod |
Типичная prod-конфигурация брокера:
listeners=SASL_SSL://0.0.0.0:9092
security.inter.broker.protocol=SASL_SSL
sasl.mechanism.inter.broker.protocol=SCRAM-SHA-512
Аутентификация (SASL)
| Механизм | Когда использовать |
|---|---|
SASL/PLAIN | Dev/test; в prod только с TLS и централизованным хранением паролей |
SASL/SCRAM-SHA-256/512 | Стандарт для prod: хеши на брокере, без plain-text пароля в конфиге клиента |
SASL/GSSAPI (Kerberos) | Корпоративный AD/LDAP, единый SSO |
SASL/OAUTHBEARER | OAuth 2.0 / OIDC (облако, zero-trust) |
Клиент задаёт механизм через sasl.mechanism и JAAS-конфиг или sasl.jaas.config.
Шифрование
- In-transit — TLS между client↔broker и broker↔broker (
ssl.keystore.location,ssl.truststore.locationна брокере; CA для клиентов). - At rest — шифрование диска на уровне ОС/облака (EBS encryption, LUKS); Kafka не шифрует log-сегменты сама.
- End-to-end — приложение шифрует payload до отправки; брокер видит только байты (для особо чувствительных данных).
Авторизация (ACL)
Управление через AclAuthorizer (или RBAC в Confluent). Пример через CLI:
kafka-acls.sh --bootstrap-server broker:9092 \
--add --allow-principal User:alice \
--operation Read --operation Describe \
--topic payments
| Операция | Смысл |
|---|---|
| READ / WRITE | Чтение / запись в топик |
| CREATE / DELETE / ALTER | Администрирование топика |
| DESCRIBE | Метаданные без чтения данных |
| ALL | Полный доступ к ресурсу |
Ресурсы — Topic, Group, Cluster, TransactionalId. Для межброкерного трафика — отдельный principal с минимальными правами.
Аудит и квоты
- Audit logs фиксируют denied/allowed операции — отправка в SIEM.
- Quotas ограничивают produce/fetch bytes/sec по principal — защита от "шумного соседа".
Добавляйте SASL_SSL параллельно с PLAINTEXT на отдельном listener, мигрируйте клиентов, затем отключайте plaintext. Ротация сертификатов — до истечения срока, с truststore на всех клиентах.
Программное управление через AdminClient
AdminClient предоставляет программный интерфейс для административных операций без использования CLI. Основные возможности:
Управление топиками
// Java: создание топика с 3 партициями и фактором репликации 2
NewTopic topic = new NewTopic("orders", 3, (short) 2);
CreateTopicsResult result = adminClient.createTopics(Collections.singleton(topic));
result.all().get();
Управление конфигурацией
# Python (confluent-kafka) — изменение параметра топика
from confluent_kafka.admin import ConfigResource, ConfigSource
config_resource = ConfigResource(ConfigResource.Type.TOPIC, "orders")
config_resource.set_config("retention.ms", "604800000")
admin.alter_configs([config_resource])
Управление группами потребителей
// C# (Confluent.Kafka): получение информации о группе
var groups = admin.ListConsumerGroupsAsync().Result;
var description = admin.DescribeConsumerGroupsAsync(new[] { "order-processors" }).Result;
AdminClient поддерживает асинхронные операции, обработку ошибок через исключения и работу с метаданными кластера (список брокеров, топиков, партиций).
Утилиты командной строки для администрирования
Стандартный дистрибутив Kafka включает скрипты в каталоге bin/:
Управление топиками
# Создание топика
kafka-topics.sh --bootstrap-server broker1:9092 --create \
--topic events --partitions 6 --replication-factor 3
# Описание топика
kafka-topics.sh --bootstrap-server broker1:9092 --describe --topic events
# Удаление топика (требует включённого delete.topic.enable)
kafka-topics.sh --bootstrap-server broker1:9092 --delete --topic deprecated-topic
Работа с потребителями
# Просмотр групп потребителей
kafka-consumer-groups.sh --bootstrap-server broker1:9092 --list
# Описание состояния группы
kafka-consumer-groups.sh --bootstrap-server broker1:9092 \
--describe --group payment-processors
# Сброс смещений (требует остановки потребителей группы)
kafka-consumer-groups.sh --bootstrap-server broker1:9092 \
--group analytics --reset-offsets --to-earliest --execute --topic clicks
Диагностика кластера
# Проверка работоспособности брокера
kafka-broker-api-versions.sh --bootstrap-server broker1:9092
# Анализ распределения партиций
kafka-topics.sh --bootstrap-server broker1:9092 --describe \
--under-replicated-partitions
В режиме KRaft утилиты используют --bootstrap-server. Параметр --zookeeper относится к legacy-кластерам Kafka 3.x.
Настройка и использование клиентских API
Producer API
Ключевые параметры конфигурации:
bootstrap.servers— начальный список брокеровacks— уровень подтверждения —0(без подтверждения),1(лидер),all(все реплики в ISR)retriesиretry.backoff.ms— политика повторных попытокcompression.type— алгоритм сжатия (none,gzip,snappy,lz4,zstd)batch.sizeиlinger.ms— параметры батчинга для повышения пропускной способности (теория — Пакетная работа с данными)enable.idempotence— идемпотентная отправка (см. exactly-once)schema.registry.url— при Avro/Protobuf через Confluent serializers
Headers — метаданные записи (trace-id, event-type); не участвуют в выборе партиции.
Interceptors — hooks до/после send для метрик и аудита без дублирования кода в каждом сервисе.
Пример отправки с ключом для упорядоченности:
ProducerRecord<String, String> record =
new ProducerRecord<>("user-events", userId, eventJson);
producer.send(record, (metadata, exception) -> {
if (exception != null) handleFailure(exception);
});
Consumer API
Ключевые параметры:
group.id— идентификатор группы потребителейauto.offset.reset— поведение при отсутствии смещения (earliest,latest)enable.auto.commit— автоматическое подтверждение смещений (рекомендуетсяfalseдля контроля)max.poll.records— максимум записей за один вызовpoll()session.timeout.msиheartbeat.interval.ms— параметры жизненного цикла потребителя в группе
Рекомендуемый цикл обработки с ручным подтверждением:
while True:
records = consumer.poll(timeout_ms=1000)
for record in records:
process(record)
consumer.commit_sync() # Подтверждение после успешной обработки
Для обработки ошибок необходимо обрабатывать исключения CommitFailedException (возникает при перебалансировке во время коммита) и реализовывать стратегию повторных попыток для идемпотентной обработки сообщений.
Установка и эксплуатация
Материал по установке вынесен в соседние статьи, чтобы не дублировать шаги:
| Задача | Статья |
|---|---|
| Первый запуск, KRaft, Docker, console producer/consumer | Apache Kafka — потоковая обработка |
| KRaft single-node на Windows, Java CRM + PostgreSQL | Java-приложение с Apache Kafka и PostgreSQL |
| CLI, параметры broker/producer/consumer, anti-patterns | Справочник по Apache Kafka |
| Мониторинг lag, URP, SLO | раздел "Мониторинг и SLO" в этой статье |
См. также
Базовая глава в разделе "Система и сеть" — Apache Kafka — потоковая обработка данных.
Как внедрять event streaming поэтапно
Чтобы снизить риск, внедряйте потоковую архитектуру итеративно:
- Выберите 1-2 доменных события с понятной ценностью для бизнеса.
- Поднимите producer + consumer + мониторинг lag/error.
- Добавьте схему события (Avro/JSON Schema) и правила совместимости в Schema Registry.
- Подключите DLQ/повторную обработку.
- Только после стабилизации расширяйте на другие домены.
Такой путь быстрее приводит к устойчивому результату, чем "большой взрыв" с полной миграцией всех интеграций сразу.