Kafka: события, команды, Topics и Keys — Часть 6

Содержание
Коротко: Kafka может работать и как очередь задач, и как поток событий. Production-дизайн начинается с выбора типа сообщения, границ Topic, Key, количества Partition, Retention и контракта события.

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

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

Начнём с назначения сообщения, а затем последовательно выберем остальные параметры.

Kafka в роли очереди задач

Допустим, после покупки нужно сформировать PDF-счёт. Сервис заказов отправляет задачу, а один из свободных обработчиков выполняет её. Такие обработчики обычно называют воркерами.

В Kafka воркеры объединяются в Consumer Group. Каждую партицию внутри группы в конкретный момент читает только один consumer. Так задачи распределяются между участниками группы:

Topic: invoice-generation

Partition 0 -> Worker A
Partition 1 -> Worker B
Partition 2 -> Worker C

Consumer Group: invoice-workers

Если задач стало больше, чем воркеры успевают обработать, они остаются в журнале. Например, во время распродажи поступило 10 000 запросов на счета, а обработчики успели выполнить только 2 000. Остальные можно прочитать позже, пока не истёк срок хранения.

После успешной обработки consumer сохраняет свою позицию чтения — делает commit offset. Для каждой назначенной партиции позиция хранится отдельно. Само сообщение при этом остаётся в Kafka.

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

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

Kafka как поток событий

Теперь представим другое сообщение: OrderPaid — «заказ оплачен». Оно сообщает о факте, который уже произошёл.

Этот факт нужен сразу нескольким сервисам:

Topic: commerce.orders.events
             |
             +--> delivery-group: подготовить доставку
             |
             +--> analytics-group: обновить статистику продаж
             |
             +--> notification-group: отправить подтверждение

Каждый сервис использует свою consumer group. Поэтому доставке, аналитике и уведомлениям доступны одни и те же события, а позиции чтения у них независимы.

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

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

Очередь задач и поток событий

Разница прежде всего в назначении сообщения.

ХарактеристикаОчередь задачПоток событий
Что передаёмПросьбу выполнить действиеФакт о произошедшем
ПримерСформировать счётЗаказ оплачен
Кто обрабатываетОжидаемый исполнитель, воркеры одной группыНесколько заинтересованных сервисов со своими группами
Что отслеживаемВыполнена ли задачаДо какой позиции каждый сервис обработал поток
Можно ли читать повторноДа, пока запись хранится в KafkaДа, пока запись хранится в Kafka

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

Один поток событий может при этом запускать фоновые задачи. Например, сервис уведомлений читает OrderPaid и отправляет письмо. Для него это работа, но исходное сообщение остаётся фактом об оплате, доступным и другим сервисам.

Сначала определите назначение сообщения

Перед созданием топика полезно договориться, что именно приложение публикует. От этого зависит, кто отвечает за обработку и какой смысл получатель должен извлечь из сообщения.

Команда

Команда просит выполнить действие: GenerateInvoice, SendEmail, ReserveInventory.

Например, сервис заказов отправляет GenerateInvoice сервису биллинга. У команды есть ожидаемый исполнитель, а результат нужно отслеживать: счёт сформирован или задача завершилась ошибкой.

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

Событие

Событие фиксирует результат: OrderCreated, PaymentCompleted, UserRegistered.

Например, сервис платежей публикует PaymentCompleted после успешной оплаты. Он отвечает за достоверность этого факта. Сервис доставки решает начать сборку заказа, аналитика учитывает продажу, а уведомления отправляют письмо.

Такое разделение уменьшает зависимость отправителя от получателей: сервису платежей не нужно знать обо всех будущих подписчиках.

Технический поток

Через Kafka также передают логи приложений, метрики, аудит и изменения строк базы данных. Последний сценарий называют CDC (Change Data Capture): изменения в базе превращаются в поток сообщений.

Здесь требования определяются назначением данных. Отладочные логи могут иметь короткий срок хранения, аудит — длительный. Для некоторых метрик допустимы пропуски, а для потока изменений базы они могут нарушить состояние системы-получателя.

Поэтому «технические сообщения» не означают автоматически «некритичные». Для каждого потока нужно определить допустимую потерю данных, срок хранения и требования к порядку.

Не создавайте один Topic для всего

Когда тип сообщения понятен, нужно выбрать, какие сообщения будут жить в одном топике.

На старте удобно создать общий events и отправлять туда всё:

UserRegistered
OrderCreated
PaymentFailed
EmailSent
ProductUpdated

Но топик имеет общие настройки хранения и доступа. Если отладочные события достаточно хранить сутки, а платёжные — месяц, внутри одного топика настроить эти сроки по типу сообщения не получится.

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

Обычно удобнее разделять потоки по бизнес-назначению:

users.events
orders.events
payments.events
notifications.commands

Так для платежей можно отдельно задать права доступа и срок хранения, а для каталога — увеличить число партиций при росте нагрузки.

Разделение не изолирует ресурсы брокеров автоматически: разные топики могут работать на одних машинах. Но оно позволяет управлять потоками независимо.

Но Topic на каждый тип события — тоже не всегда хорошо

Можно пойти в другую крайность и создать отдельный топик для каждого изменения заказа:

order-created
order-paid
order-cancelled
order-shipped
order-delivered

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

Если у событий одинаковые требования к доступу и хранению, часто удобнее объединить их в orders.events, использовать orderId как ключ и указывать тип внутри сообщения:

{
  "eventType": "OrderPaid",
  "eventId": "5c72a8f1",
  "orderId": "48291",
  "occurredAt": "2026-08-06T09:30:00Z",
  "data": {
    "amountMinor": 1599000,
    "currency": "RUB"
  }
}

Если же один тип событий содержит чувствительные данные или требует другого срока хранения, отдельный топик может быть оправдан.

Граница топика должна следовать требованиям потока. Само количество типов сообщений не даёт готового ответа.

Как именовать Topics

Имя топика должно помогать понять, какие данные он содержит. Например, можно принять соглашение:

<domain>.<entity>.<message-type>

Тогда имена будут выглядеть так:

commerce.orders.events
commerce.payments.events
identity.users.events
notifications.email.commands

commerce обозначает область системы, orders — объект, а events — назначение сообщений. Это соглашение команды, а не обязательный формат Kafka.

Главное — применять его последовательно. Смесь payment_events_v2, UserUpdateTopic и email-queue-prod-new затрудняет поиск и со временем делает имена неоднозначными.

Как выбрать Key

После выбора топика нужно решить, как распределять его сообщения между партициями.

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

Это позволяет сохранять порядок для конкретного объекта и распределять разные объекты между обработчиками.

Пример с заказами

Для событий заказа естественным ключом будет orderId. Ключ передаётся как отдельная часть записи Kafka, а не просто поле внутри JSON:

Key: order-48291
Value: {"eventType": "OrderPaid", "orderId": "48291"}

Если события публикуются последовательно, журнал одной партиции будет содержать:

OrderCreated -> PaymentCompleted -> DeliveryStarted -> OrderDelivered

Другие заказы с другими ключами могут попасть в другие партиции и обрабатываться параллельно.

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

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

Пример с пользователями

Для изменений профиля можно использовать userId. Тогда EmailChanged, PasswordChanged и UserBlocked одного пользователя попадут в одну партицию.

Выбирать ключ нужно по объекту, для которого важен порядок. Если требуется последовательность действий внутри заказа, используйте идентификатор заказа; если внутри профиля — идентификатор пользователя.

Плохой ключ: константа

Если у всех сообщений ключ orders, они попадут в одну партицию:

Partition 0: 100% нагрузки
Partition 1:   0% нагрузки
Partition 2:   0% нагрузки

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

Плохой ключ: поле с малым количеством значений

Допустим, ключ — статус заказа: created, paid или cancelled.

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

Такой ключ одновременно ограничивает распределение нагрузки и не обеспечивает нужный порядок.

Плохой ключ: текущее время

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

OrderCreated     -> Partition 1
PaymentCompleted -> Partition 4
DeliveryStarted  -> Partition 0

Для аналитики, которой порядок отдельных заказов не важен, это может быть допустимо. Для последовательного изменения состояния заказа такой ключ не подходит.

Если порядок вообще не требуется, можно отправлять сообщения без ключа и использовать распределение producer. Придумывать искусственный ключ только ради наличия поля не нужно.

Как избежать Hot partition

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

Все его события попадают в одну партицию. Другие consumers свободны, но внутри группы они не могут разделить чтение этой партиции между собой.

Сначала нужно уточнить, действительно ли требуется общий порядок всех событий продавца.

Выбрать более точный объект

Если порядок нужен только внутри заказа, вместо sellerId лучше использовать orderId.

Заказы крупного продавца смогут распределяться по разным партициям, а события каждого заказа останутся вместе. Это часто позволяет убрать перекос без усложнения обработки.

Использовать составной ключ

Если идентификатор продавца нужно сохранить в ключе, можно добавить номер группы — shard:

seller-1-0
seller-1-1
seller-1-2

Например, номер shard вычисляется из orderId. Тогда события одного заказа получают одинаковый составной ключ, а разные заказы распределяются между несколькими ключами.

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

Разделить поток

Особенно крупного продавца можно вынести в отдельный топик с собственными настройками.

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

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

Сколько Partition создавать

Ключ определяет распределение сообщений, а число партиций — доступный параллелизм чтения внутри одной consumer group.

При четырёх партициях одновременно читать их смогут максимум четыре участника группы:

Partition 0 -> Consumer A
Partition 1 -> Consumer B
Partition 2 -> Consumer C
Partition 3 -> Consumer D

Пятый consumer останется без назначенных партиций. Другие группы могут независимо читать этот же топик своими consumers.

Число партиций выбирают с учётом скорости записи, обработки, размера сообщений, ресурсов брокеров и роста нагрузки.

Простой подход к расчёту

Предположим, нагрузочный тест показал: один consumer с реальной бизнес-логикой обрабатывает 5 000 сообщений в секунду. Ожидаемый поток — 40 000 сообщений в секунду.

40 000 / 5 000 = 8 consumers

Для восьми consumers одной группы нужно хотя бы восемь партиций. Но это расчёт без запаса: если поток вырастет или один обработчик остановится, группа начнёт отставать.

Можно проверить вариант с 12 или 16 партициями и достаточным числом consumers. Сами по себе дополнительные партиции производительность не создают: брокеры, сеть и база приложения должны выдержать суммарную нагрузку.

Расчёт также предполагает достаточно равномерное распределение. Горячая партиция может отставать даже при низкой общей загрузке.

Почему нельзя просто создать тысячу Partition

У каждой реплики партиции есть файлы журнала, индексы и служебное состояние. Например, 1 000 партиций с replication.factor=3 означают 3 000 реплик, распределённых по брокерам.

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

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

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

Планируйте рост заранее

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

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

partition = hash(key) % partition_count

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

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

Retention: сколько хранить сообщения

Когда распределение нагрузки выбрано, нужно решить, сколько времени сообщения должны оставаться доступными. Политику хранения называют Retention.

Это напрямую влияет на восстановление. Если consumer не работал три дня, а сообщения хранятся только сутки, часть пропущенных событий уже будет удалена. Возврат к старому offset не вернёт эти данные.

Нужно учитывать не только время простоя, но и время, за которое consumer догонит поток после запуска. Ему придётся обрабатывать накопившиеся сообщения одновременно с новыми.

Для выбора срока хранения полезно выяснить:

  • как долго допустим простой consumers;
  • за какое время они смогут разобрать накопившиеся данные;
  • нужно ли пересчитывать историю;
  • можно ли восстановить события из другого источника;
  • сколько места доступно и какие ограничения действуют на хранение данных.

Хранение по времени

Для удаления старых записей используют cleanup.policy=delete. Например, семь дней хранения задаются так:

cleanup.policy=delete
retention.ms=604800000

Kafka удаляет данные целыми файлами-сегментами. Поэтому сообщение не обязательно исчезнет ровно через семь дней: удаление зависит от возраста сегмента и фоновой проверки.

Такой срок может подойти для потока просмотров товаров, если аналитика обычно читает его сразу, а команде нужна возможность пересчитать статистику за последнюю неделю.

Хранение по размеру

retention.bytes ограничивает объём журнала одной партиции. Если для топика с десятью партициями задан лимит 10 GiB, ориентир для одной копии всего топика — 100 GiB.

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

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

Log Compaction

Для некоторых потоков важна не вся история, а последнее известное состояние каждого объекта. Например, сервису нужен локальный кэш профилей пользователей.

Для этого можно использовать cleanup.policy=compact — компакцию журнала. Она работает по ключу сообщения:

Offset 0: user-1 -> name=Ivan
Offset 1: user-2 -> name=Anna
Offset 2: user-1 -> name=Sergey

После фоновой очистки старое значение user-1 может быть удалено:

Offset 1: user-2 -> name=Anna
Offset 2: user-1 -> name=Sergey

Порядок оставшихся записей и их offsets сохраняются. Компакция не происходит мгновенно: consumer может увидеть несколько версий одного ключа и должен применять их последовательно.

Для удаления объекта отправляют запись с его ключом и значением null. Её называют tombstone: для consumer это сигнал удалить объект из своего состояния. Такая запись тоже хранится ограниченное время, поэтому правила восстановления кэша нужно учитывать отдельно.

Компакция подходит для полного состояния профилей, конфигураций и некоторых CDC-потоков. Если сообщение содержит только изменение вроде «прибавить 100 бонусов», сохранение последней записи не восстановит итоговый баланс. В таком случае нужна история изменений или отдельный поток полного состояния. Описание компакции Kafka.

Проектируйте формат события

Теперь у нас есть топик, ключ, партиции и политика хранения. Осталось договориться, что именно находится внутри сообщения.

Kafka хранит ключ и значение как байты. Она сама не знает, что поле total означает стоимость заказа. Это знание должно быть частью контракта между отправителем и получателями.

Например:

{
  "eventId": "0ae2342d-105d-45ad-8d91-441b7115afbc",
  "eventType": "OrderCreated",
  "eventVersion": 1,
  "occurredAt": "2026-08-06T09:30:00Z",
  "producer": "order-service",
  "data": {
    "orderId": "48291",
    "userId": "1732",
    "amountMinor": 1599000,
    "currency": "RUB"
  }
}

Здесь eventType описывает факт, occurredAt — время, когда он произошёл, а producer — сервис-источник. amountMinor хранит сумму в минимальных единицах валюты: для RUB это копейки. Единицы измерения нужно явно описать в контракте.

occurredAt не обязательно совпадает с timestamp записи Kafka. Например, событие произошло в 12:00, а из-за временной недоступности Kafka было отправлено в 12:05. Отдельное поле позволяет сохранить бизнес-время независимо от настроек временной метки Kafka.

Зачем нужен Event ID

При повторной доставке consumer может получить событие ещё раз. eventId позволяет понять, что это тот же факт, а не новая операция.

Например, сервис бонусов получает PaymentCompleted и начисляет баллы. Перед повторным начислением он проверяет идентификатор события.

Идентификатор нужно сохранять при повторных отправках того же события. Если при каждом retry генерировать новый eventId, получатель не сможет распознать дубль.

Кроме того, проверка идентификатора и изменение баланса должны быть согласованы, например выполнены в одной транзакции базы. Само наличие поля в JSON не делает обработку идемпотентной.

Зачем нужна версия события

Со временем формат меняется. Например, простое поле суммы заменяется объектом:

{
  "amount": {
    "valueMinor": 1599000,
    "currency": "RUB"
  }
}

Старый consumer может ожидать amountMinor и не понимать новый формат. eventVersion помогает явно выбрать нужный обработчик или преобразование.

Но увеличение номера версии не исправляет несовместимость автоматически. Сначала нужно подготовить consumers к новому формату, а затем начать его публиковать. При повторном чтении истории им могут встретиться обе версии.

Не меняйте смысл существующего поля

Допустим, price хранит рубли:

{"price": 1500}

После обновления producer начинает передавать в том же поле копейки:

{"price": 150000}

Тип поля не изменился, поэтому consumer может продолжить работу и посчитать цену в сто раз больше.

Такую ошибку сложнее заметить, чем отказ разбора JSON. Лучше ввести новое поле с явным смыслом, например amountMinor, и согласовать переход с получателями.

Предпочитайте совместимые изменения

Предположим, в событие заказа добавляется deliveryComment. Если старые consumers игнорируют неизвестные поля, а новые допускают отсутствие комментария в старых событиях, сервисы можно обновлять постепенно.

Совместимость нужно проверять в обе стороны:

  • понимает ли новый consumer старые сообщения, оставшиеся в Kafka;
  • понимает ли старый consumer новые сообщения, которые уже публикует producer.

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

Конкретные правила зависят от формата и реализации читателей: изменение JSON, Avro или Protobuf нельзя оценивать только по внешнему виду.

Когда нужен Schema Registry

Если один поток читают несколько сервисов, устной договорённости о полях становится мало. Нужна машиночитаемая схема: например, Avro, Protobuf или JSON Schema.

Schema Registry — отдельный сервис, который хранит схемы и проверяет их совместимость при регистрации новых версий. Он не является встроенным валидатором всех записей Kafka: приложения должны использовать соответствующие средства сериализации и проверки.

Основные режимы совместимости:

  • BACKWARD: новый consumer с новой схемой может читать данные предыдущей схемы;
  • FORWARD: consumer с предыдущей схемой может читать данные новой схемы;
  • FULL: выполняются оба требования.

Например, при BACKWARD сначала можно обновить consumers, сохранив возможность читать старые сообщения, а затем producers. Если требуется перечитывать всю историю с несколькими версиями схем, нужна проверка совместимости со всеми предыдущими версиями, например BACKWARD_TRANSITIVE.

Схема помогает обнаруживать структурные несовместимости, но не заменяет описание смысла полей. Смена рублей на копейки может пройти проверку типов и всё равно сломать бизнес-логику. Подробнее о совместимости схем.

Документируйте владельца Topic

Даже хорошо спроектированный поток трудно сопровождать, если неизвестно, кто за него отвечает.

Для каждого топика полезно сохранить небольшую карточку:

Topic: commerce.orders.events
Owner: команда заказов
Producer: order-service
Consumers: доставка, уведомления, аналитика
Message type: события жизненного цикла заказа
Key: orderId
Partitions: 12
Replication factor: 3
Retention: 14 дней, cleanup.policy=delete
Schema: OrderEvent, проверка совместимости в CI
Ordering: порядок записи внутри партиции, ключом служит orderId

Числа в примере условные: они должны следовать из нагрузки и требований восстановления.

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

Контрольный список перед запуском

Перед созданием топика проверьте, что команда может ответить на основные вопросы.

ОбластьЧто выяснить
НазначениеЭто команда, событие или технический поток? Кто публикует и кто обрабатывает?
ПорядокДля какого объекта он нужен? Какой ключ используем? Возможны ли горячие ключи?
НагрузкаСколько данных поступает? Как быстро работают consumers? Какой запас ресурсов нужен?
ХранениеСколько длится допустимый простой? Успеют ли consumers догнать поток до удаления данных?
КонтрактЕсть ли стабильный eventId, понятные единицы измерения и проверка совместимости?
ЭксплуатацияКто владеет топиком? Кто реагирует на ошибки и рост отставания consumers?

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

Итог

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

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

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

Что дальше?

Мы определили, какие данные передаём и как организуем поток. В следующей части разберём, что происходит при сбоях: producer не получил подтверждение, consumer выполнил операцию и упал, а сообщение снова пришло на обработку.

На этих примерах рассмотрим acks, идемпотентность, сохранение offset, retry, DLQ, Transactional Outbox и мониторинг отставания consumers.