Kafka: жизненный цикл сообщения от Producer до Consumer — Часть 7
Мы уже разобрали отдельные компоненты Kafka:
Producer;Topic;Partition;Broker;LeaderиFollower;Consumer Group;Offset.
Теперь соберём всё вместе и посмотрим, что происходит с сообщением на каждом этапе. Это тот момент, когда отдельные детали наконец перестают лежать кучей на столе и начинают напоминать систему.
Упрощённо полный путь выглядит так:
Producer
|
v
Выбор Partition
|
v
Leader Broker
|
v
Append-Only Log
|
v
Follower Replicas
|
v
Consumer
|
v
Offset
В исходном материале этот процесс описан короткой последовательностью:
Producer -> Topic Partition -> Broker -> Consumer
Но за этой простой схемой скрывается несколько важных шагов. Как обычно в инфраструктуре: стрелочка на диаграмме выглядит невинно, а внутри у неё маленькая жизнь.
Шаг 1. Producer создаёт сообщение
Представим сервис заказов.
После успешного оформления покупки он создаёт событие:
{
"orderId": 48291,
"userId": 1732,
"status": "created",
"total": 15990
}
Само сообщение может содержать:
Key
Value
Timestamp
Headers
Например:
Key: 48291
Value:
{
"orderId": 48291,
"status": "created"
}
Timestamp:
2026-08-06T09:30:00Z
Headers:
trace-id=fa72c8
content-type=application/json
Kafka рассматривает содержимое сообщения как набор байтов. Она не знает, что внутри находится заказ, платёж или пользователь. Смысл данных определяют Producer и Consumer.
Шаг 2. Producer определяет Topic
Producer должен знать, в какой Topic отправить сообщение.
Например:
orders
В коде операция концептуально выглядит так:
send(
topic = "orders",
key = "48291",
value = order
)
Topic является логическим именем потока.
Producer при этом не обязан заранее знать:
- на каком
BrokerхранитсяTopic; - сколько у него
Partition; - какой
Brokerсейчас являетсяLeader; - где расположены резервные реплики.
Эту информацию клиент получает из метаданных Kafka-кластера.
Шаг 3. Выбирается Partition
Следующий вопрос:
В какую Partition попадёт сообщение?
Если у сообщения есть Key, Producer использует его для выбора Partition.
Упрощённая модель:
hash(key) % partition_count
Допустим:
Key = 48291
Количество Partition = 3
После вычисления Producer получает:
Partition 1
Все последующие сообщения с тем же ключом при неизменном количестве Partition также будут направлены в эту Partition.
Это позволяет сохранить порядок событий одного заказа:
OrderCreated
PaymentStarted
PaymentCompleted
DeliveryStarted
OrderDelivered
Исходный материал подчёркивает, что одинаковый ключ направляет связанные сообщения в одну Partition и тем самым обеспечивает упорядоченную обработку внутри неё.
Что происходит без Key
Если Key отсутствует, Producer старается равномерно распределить сообщения между Partition.
Упрощённо это может выглядеть так:
Message 1 -> Partition 0
Message 2 -> Partition 1
Message 3 -> Partition 2
Message 4 -> Partition 0
В современных клиентах часто применяется так называемое sticky partitioning: некоторое количество сообщений временно направляется в одну Partition, чтобы эффективнее формировать Batch.
Для архитектурного понимания важно следующее:
Без ключа Kafka распределяет нагрузку, но не гарантирует порядок логически связанных сообщений.
Шаг 4. Producer находит Leader
После выбора Partition необходимо определить Broker, который принимает запись.
Представим:
Topic: orders
Partition 0 -> Broker 2
Partition 1 -> Broker 3
Partition 2 -> Broker 1
Для выбранной Partition 1 Leader находится на Broker 3.
Producer получает эту информацию из метаданных кластера и устанавливает соединение непосредственно с нужным Broker.
Producer
|
v
Broker 3
Leader of orders-1
Producer не отправляет сообщение случайному Broker и не рассылает его всем узлам кластера. Запись направляется конкретному Leader выбранной Partition.
Что произойдёт, если метаданные устарели
Представим, что Producer считает лидером Broker 3, но тот уже вышел из строя.
Kafka могла выбрать нового Leader:
Старый Leader: Broker 3
Новый Leader: Broker 1
Producer получает ошибку, обновляет метаданные и повторяет отправку уже на актуальный Broker.
Producer
|
| ошибка
v
Broker 3
Producer
|
| обновление metadata
v
Broker 1
Такие переключения обычно обрабатываются клиентской библиотекой автоматически.
Шаг 5. Сообщение попадает в Batch
Producer обычно не отправляет каждое сообщение отдельным сетевым запросом.
Он временно накапливает несколько Record:
Record 1
Record 2
Record 3
Record 4
После чего отправляет их одной пачкой:
Batch
Вместо четырёх сетевых операций выполняется одна.
4 сообщения
|
v
1 сетевой запрос
Батчинг уменьшает:
- количество сетевых запросов;
- число системных вызовов;
- нагрузку на CPU;
- накладные расходы
Broker.
Исходная статья выделяет batching как одну из основных причин высокой производительности Kafka.
Шаг 6. Leader дописывает сообщение в журнал
Leader получает Batch и записывает сообщения в журнал Partition.
Partition 1
Offset 1001
Offset 1002
Offset 1003
Offset 1004
Каждое новое сообщение добавляется строго в конец.
Старые записи
Offset 1001
Offset 1002
Offset 1003
Новое сообщение
Offset 1004
Kafka не вставляет запись в середину и не изменяет предыдущие сообщения.
Именно поэтому журнал называется append-only log.
Шаг 7. Сообщению назначается Offset
Каждый Record получает последовательный Offset внутри Partition.
Например:
Topic: orders
Partition: 1
Offset: 1004
Полная позиция сообщения определяется не одним Offset, а сочетанием:
Topic + Partition + Offset
Например:
orders / 1 / 1004
Offset 1004 может одновременно существовать в другой Partition:
orders / 0 / 1004
orders / 1 / 1004
orders / 2 / 1004
Это разные сообщения.
Offset является локальным порядковым номером внутри конкретной Partition.
Шаг 8. Followers копируют сообщение
После записи на Leader резервные реплики начинают копировать новые данные.
Broker 3
Leader
Offset 1004
|
+----------------+
| |
v v
Broker 1 Broker 2
Follower Follower
Offset 1004 Offset 1004
Followers самостоятельно запрашивают новые записи у Leader.
Как только они догоняют его, реплики считаются синхронизированными.
Kafka использует модель Leader-Follower:
Leaderобслуживает операции чтения и записи;Followersподдерживают копии журнала;Controllerследит заBrokerи выбирает новогоLeaderпри отказе текущего.
Шаг 9. Producer получает подтверждение
Момент получения подтверждения зависит от настроек надёжности.
Концептуально возможны три модели.
Не ждать подтверждения
Producer -> Kafka
Producer продолжает работу
Это быстро, но Producer не знает, было ли сообщение принято.
Ждать Leader
Producer -> Leader
Leader записал сообщение
Leader -> OK
Producer получает подтверждение после записи на основной Broker.
Ждать синхронные реплики
Producer -> Leader
Leader -> Followers
Реплики синхронизированы
Leader -> OK
Это надёжнее, но может увеличить задержку.
Исходный материал описывает сам механизм Leader-Follower и репликации, но не разбирает конфигурацию подтверждений подробно. Поэтому конкретные настройки Producer здесь показаны как дополнительное архитектурное пояснение, а не как перевод исходного текста.
Шаг 10. Consumer запрашивает данные
Kafka использует Pull-модель.
Это означает, что Broker не пытается самостоятельно отправить сообщение Consumer.
Вместо этого Consumer спрашивает:
Есть ли новые сообщения после Offset 1003?
Broker отвечает:
Да:
Offset 1004
Offset 1005
Offset 1006
Полный поток выглядит так:
Consumer
| fetch request
v
Broker
| records
v
Consumer
Такой подход даёт Consumer контроль над скоростью обработки.
Почему Kafka выбрала Pull-модель
Контроль нагрузки
Consumer забирает столько данных, сколько способен обработать.
Быстрый Producer
|
v
Kafka
|
v
Медленный Consumer
Broker не перегружает Consumer принудительной отправкой сообщений.
Данные остаются в Topic, пока приложение постепенно их обрабатывает.
Эффективное чтение Batch
Consumer может запросить не одно сообщение, а сразу большой блок.
Fetch request
|
v
1000 сообщений
Это уменьшает количество сетевых обращений и повышает пропускную способность.
Возможность вернуться назад
Поскольку прочитанные сообщения остаются в журнале, Consumer может изменить свою позицию.
Текущий Offset: 10000
Новый Offset: 5000
После этого он повторно прочитает старые события.
Именно возможность сбросить Offset делает Kafka пригодной для:
- повторной обработки;
- восстановления данных;
- пересчёта аналитики;
- создания нового
Consumerна основе старой истории.
Исходная статья выделяет управление нагрузкой, batching и возможность перемотки как три главных преимущества Pull-модели.
Шаг 11. Kafka распределяет Partition между Consumer
Представим группу:
Consumer Group: order-processing
В ней три Consumer:
Consumer A
Consumer B
Consumer C
А в Topic три Partition:
Partition 0
Partition 1
Partition 2
Kafka распределяет их:
Partition 0 -> Consumer A
Partition 1 -> Consumer B
Partition 2 -> Consumer C
Внутри одной группы одну Partition в конкретный момент читает только один Consumer.
Это позволяет:
- распределять нагрузку;
- сохранять порядок внутри
Partition; - избегать параллельной обработки одного журнала несколькими участниками группы.
Шаг 12. Consumer обрабатывает сообщение
Consumer получает событие:
{
"orderId": 48291,
"status": "created"
}
После этого выполняет бизнес-логику.
Например:
Проверить заказ
Зарезервировать товар
Создать запись о доставке
Отправить результат
Kafka не знает, успешно ли выполнилась эта логика.
Для Broker получение сообщения Consumer и его фактическая обработка приложением — разные вещи.
Именно здесь начинается зона ответственности разработчика.
Шаг 13. Consumer фиксирует позицию
После обработки Consumer сохраняет свою позицию чтения.
Например:
Topic: orders
Partition: 1
Следующий Offset: 1005
Это означает:
Все сообщения до Offset 1004 обработаны, продолжать нужно с 1005.
Kafka сохраняет Offset отдельно для каждой Consumer Group.
Group: delivery
orders / Partition 1 / Offset 1005
Другая группа может находиться в совершенно другой позиции:
Group: analytics
orders / Partition 1 / Offset 850
Они не мешают друг другу.
Что произойдёт после перезапуска Consumer
Допустим, Consumer выключился после обработки Offset 1004.
После запуска он читает сохранённую позицию:
Next Offset = 1005
И продолжает работу:
Offset 1005
Offset 1006
Offset 1007
Благодаря этому приложение не обязано хранить всю очередь локально.
Kafka запоминает его положение в журнале.
Что произойдёт при падении Consumer
Представим:
Partition 1 -> Consumer B
Consumer B перестал отвечать.
Kafka запускает перераспределение Partition:
Partition 1 -> Consumer A
Новый Consumer начинает чтение с последнего сохранённого Offset группы.
Последний сохранённый Offset: 1005
В зависимости от того, когда приложение сохранило Offset относительно бизнес-операции, часть сообщений может быть обработана повторно или пропущена.
Где могут появиться повторные сообщения
Рассмотрим ситуацию:
1. Consumer получил сообщение
2. Выполнил операцию в базе
3. Упал до сохранения Offset
После перезапуска Kafka видит старый Offset и снова отдаёт то же сообщение.
Одно событие
|
v
Обработано дважды
Поэтому Consumer обычно проектируют идемпотентным.
Это означает:
Повторная обработка одного сообщения не должна приводить к неправильному результату.
Например, вместо безусловной вставки:
INSERT INTO payments (...)
можно использовать уникальный идентификатор операции и проверять, не был ли платёж обработан ранее.
Где сообщение может быть пропущено
Теперь другой порядок:
1. Consumer получил сообщение
2. Сохранил Offset
3. Упал до выполнения бизнес-операции
После запуска Kafka продолжит чтение со следующего Offset.
Старое сообщение повторно не придёт, хотя бизнес-логика не была завершена.
Поэтому момент фиксации Offset критически важен.
Три модели обработки
At-most-once
Сообщение обрабатывается не более одного раза.
Позиция сохраняется до выполнения бизнес-логики.
Commit Offset
|
v
Process Message
Плюс:
- нет повторной обработки.
Минус:
- при сбое сообщение может быть потеряно для приложения.
At-least-once
Сообщение обрабатывается как минимум один раз.
Сначала выполняется бизнес-логика, затем сохраняется Offset.
Process Message
|
v
Commit Offset
Плюс:
- сообщение не пропускается из-за раннего commit.
Минус:
- возможны дубликаты.
На практике это одна из самых распространённых моделей, если Consumer является идемпотентным.
Exactly-once
Сообщение оказывает эффект ровно один раз.
Это самая сложная модель, особенно если обработка затрагивает внешние системы:
Kafka
Database
HTTP API
Payment Provider
Гарантия Kafka сама по себе не может автоматически сделать произвольную внешнюю операцию «ровно один раз». Kafka не знает, что такое «списать деньги красиво и без дублей» — это уже зона ответственности приложения.
Для этого могут потребоваться:
- транзакции;
- идемпотентные операции;
- уникальные ключи;
- Inbox/Outbox-паттерны;
- согласованная фиксация результата и
Offset.
Важно: at-most-once, at-least-once и exactly-once не рассматриваются в загруженной исходной статье. Этот раздел добавлен как отдельное пояснение, чтобы завершить обещанный разбор жизненного цикла сообщения.
Полный путь сообщения
Соберём всё в одну последовательность.
1. Producer создаёт Record
2. Указывает Topic
3. Partitioner выбирает Partition
4. Producer находит Leader Broker
5. Сообщение попадает в Batch
6. Leader дописывает Record в журнал
7. Record получает Offset
8. Followers копируют данные
9. Producer получает подтверждение
10. Consumer отправляет Fetch-запрос
11. Kafka назначает Partition Consumer
12. Consumer выполняет бизнес-логику
13. Consumer фиксирует Offset
В виде одной схемы:
Producer
|
| Record
v
Partitioner
|
| Partition 1
v
Leader Broker
|
| Append + Offset
v
Partition Log
|
+-----------> Follower 1
|
+-----------> Follower 2
|
v
Consumer Group
|
v
Consumer
|
| Business Logic
v
Offset Commit
Что важно запомнить
Kafka отвечает за:
- хранение сообщения;
- порядок внутри
Partition; - репликацию;
- доставку данных
Consumer; - хранение позиции
Consumer Group.
Приложение отвечает за:
- корректную бизнес-логику;
- момент сохранения
Offset; - обработку повторных сообщений;
- идемпотентность;
- работу с внешними системами.
Kafka может надёжно доставить событие, но она не знает, что означает «успешно провести платёж», «создать заказ» или «отправить письмо».
Эта граница между инфраструктурой и бизнес-логикой — одна из важнейших вещей при проектировании Kafka-систем. Если её забыть, потом Kafka внезапно становится виновата и в дублях, и в платежах, и почти в плохой погоде.
Что дальше?
Теперь мы прошли полный путь сообщения — от формирования Record до фиксации Offset.
В следующей части разберём практическое проектирование Kafka: тот самый слой, где хорошие идеи начинают встречаться с реальной нагрузкой, бюджетом и людьми, которые будут это сопровождать.
- как выбирать количество
Partition; - какой
Keyиспользовать; - как избежать горячих
Partition; - как именовать
Topics; - как задавать
Retention; - почему слишком большое количество
Partitionтоже опасно; - какие ошибки чаще всего допускают при внедрении Kafka.