Kafka: надёжность, retry, DLQ и эксплуатация — Часть 7

Содержание
Коротко: Надёжность Kafka строится не одной настройкой, а цепочкой решений: producer должен безопасно записывать события, consumer — корректно коммитить offset, а ошибки — уходить в retry или DLQ без блокировки партиции.

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

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

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

Гарантии доставки: что именно гарантирует Kafka

Фраза «сообщение доставлено» неоднозначна. Она может означать, что Kafka сохранила запись, consumer получил её или приложение выполнило бизнес-операцию. Это разные этапы.

Producer -> запись в Kafka -> чтение Consumer -> изменение бизнес-данных

Например, Kafka подтвердила событие об оплате, но сервис склада ещё не прочитал его. Или прочитал, но не смог зарезервировать товар из-за ошибки базы.

Параметр acks определяет условия подтверждения записи в Kafka. Репликация помогает пережить отказ брокера, а повторы producer позволяют справиться с временной ошибкой отправки.

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

Producer: acks и репликация

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

Topic:
  replication.factor=3
  min.insync.replicas=2

Producer:
  acks=all
  enable.idempotence=true

replication.factor=3 означает, что каждая партиция хранится в трёх копиях на разных брокерах. Одна копия — лидер, остальные получают новые записи от него.

acks=all требует репликации записи на все реплики текущего ISR. Напомним: это набор реплик, которые остаются достаточно синхронизированными с лидером; сам лидер тоже входит в него.

min.insync.replicas=2 задаёт минимальный размер этого набора для успешной записи с acks=all.

Например, лидер находится на брокере A, а followers — на B и C:

СостояниеЧто происходит с записью
A, B и C входят в ISRПодтверждение приходит после записи на все три реплики
C недоступен и исключён из ISRA и B могут продолжать принимать и подтверждать записи
В ISR остался только AKafka возвращает ошибку: двух актуальных реплик нет

Важно: min.insync.replicas=2 не означает «дождаться любых двух из трёх». Пока все три входят в ISR, acks=all ждёт все три. Минимум в две реплики позволяет продолжить запись после исключения одной из набора. Описание настройки Kafka.

Эта связка повышает устойчивость к отказу отдельного брокера, но не защищает от любого сочетания аварий. Кроме того, acks=all само по себе не требует немедленного fsync на накопителях: данные могут ещё находиться в файловом кэше ОС. Запись и сброс кэша.

Что происходит, если подтверждение не пришло

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

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

enable.idempotence=true позволяет Kafka распознавать повторы, которые producer выполняет автоматически, и не записывать их как новые сообщения. Этот режим требует совместимых настроек, в частности acks=all и включённых retries.

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

В Java producer отправка обычно асинхронная: успешный возврат из send() ещё не означает успешную запись в Kafka. Приложение должно проверять результат через callback или future и обрабатывать окончательную ошибку.

Время автоматических попыток ограничено, в частности, delivery.timeout.ms. Если отправка так и не завершилась, событие нужно сохранить для последующей публикации, а ошибку — показать в мониторинге. Настройки producer.

Transactional Outbox: как связать базу приложения и Kafka

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

Представим, что сервис заказов сначала сохраняет заказ в PostgreSQL, а затем публикует OrderCreated:

1. Заказ сохранён в базе
2. Приложение упало до отправки события
3. Склад и уведомления не узнали о заказе

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

Для таких случаев используют Transactional Outbox. Приложение в одной транзакции базы сохраняет и заказ, и событие в отдельную таблицу outbox:

Одна транзакция PostgreSQL:
  INSERT в orders
  INSERT в outbox_events
  COMMIT

Либо сохраняются обе записи, либо ни одна. Событие получает eventId уже здесь, и этот идентификатор сохраняется при всех попытках публикации.

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

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

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

Consumer: когда сохранять Offset

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

Эту операцию называют commit offset. Для каждой партиции сохраняется offset следующей записи, которую нужно прочитать. Например, после успешной обработки записи 100 можно сохранить позицию 101.

Offset — это позиция группы, а не удаление сообщения и не подтверждение бизнес-операции самой Kafka.

At-most-once: сначала Offset, потом обработка

Представим сервис уведомлений:

1. Consumer получил событие с offset 100
2. Сохранил следующую позицию: 101
3. Начал отправлять письмо
4. Упал до завершения отправки

После перезапуска группа продолжит с позиции 101. Событие 100 может оставаться в Kafka, но обычное восстановление группы его пропустит.

Такой порядок называют at-most-once: допускается пропуск обработки. Для обязательного уведомления об оплате это обычно неподходящее поведение.

At-least-once: сначала обработка, потом Offset

Теперь поменяем последовательность:

1. Consumer получил событие с offset 100
2. Успешно отправил письмо
3. Упал до сохранения позиции 101
4. После перезапуска снова получил событие 100

Обработка повторится. Это модель at-least-once: при восстановлении возможны повторные попытки обработки.

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

Чтобы приложение само управляло моментом commit, в Java consumer используют enable.auto.commit=false. Одного флага недостаточно: код должен сохранять позицию только после успешной обработки соответствующих записей. Документация consumer.

Где применяется Exactly-once

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

Например, обработчик прочитал оплату и сформировал событие для аналитики. В одной транзакции он публикует результат и сохраняет входную позицию. Получатели с isolation.level=read_committed читают только записи завершённых транзакций.

При корректной организации такой обработки повтор после сбоя не создаёт второй видимый результат в выходном потоке. Но обращение к внешнему платёжному API или изменение PostgreSQL не становится частью транзакции Kafka автоматически. Для этих действий нужны собственные механизмы согласования и защиты от повторов. API транзакционного producer.

Consumer должен быть идемпотентным

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

Допустим, за оплату заказа нужно начислить 100 бонусов. Consumer обновил баланс и упал до commit offset. После запуска он снова получает событие оплаты.

Если просто повторить прибавление, покупатель получит 200 бонусов. Чтобы этого не произошло, сервис хранит идентификаторы обработанных событий.

Для PostgreSQL последовательность может быть такой:

1. Начать транзакцию базы
2. Попробовать добавить eventId в processed_events
3. Если eventId новый, начислить бонусы
4. Зафиксировать транзакцию
5. Только после этого сохранить offset в Kafka

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

INSERT INTO processed_events (event_id)
VALUES ('0ae2342d-105d-45ad-8d91-441b7115afbc')
ON CONFLICT (event_id) DO NOTHING
RETURNING event_id;

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

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

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

Что делать, если сообщение не удалось обработать

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

На практике полезно различать:

  • Временные ошибки: недоступна база, внешний сервис вернул временную ошибку, истёк сетевой таймаут.
  • Ошибки данных или контракта: сообщение невозможно разобрать, поле отсутствует, версия формата неизвестна.

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

Способ восстановления выбирается по смыслу ошибки и требованиям к порядку.

Почему бесконечный Retry блокирует поток

Допустим, consumer последовательно обрабатывает записи партиции:

Offset 100: обработан
Offset 101: ошибка
Offset 102: ждёт
Offset 103: ждёт

Если запись 101 постоянно повторяется, следующие не обрабатываются. Это называют Head-of-Line Blocking: ошибка в начале очереди задерживает остальные записи.

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

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

Retry Topics: повторить позже

Если сервис уведомлений временно не может отправить письмо, событие можно перенести в отдельный retry-топик. Основной consumer продолжит работу, а другой обработчик повторит попытку позже.

orders.events
      |
      | временная ошибка отправки письма
      v
notifications.retry
      |
      | повторная попытка после задержки
      v
отправка письма

В сообщение добавляют число попыток, время следующей попытки и описание ошибки. Например, интервалы могут составлять 10 секунд, 30 секунд и 2 минуты. Это иллюстрация, а не универсальные значения.

Само имя retry-10s не заставляет Kafka ждать десять секунд. Время повтора контролирует приложение или используемый фреймворк. Чтобы восстановившийся API не получил сразу весь накопленный поток, полезно также ограничивать скорость повторов.

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

Такую схему нужно выбирать осознанно. Описание retry-топиков в Spring Kafka.

DLQ: сохранить сообщение для разбора

Dead Letter Queue, или DLQ, в Kafka обычно представляет собой отдельный топик для записей, которые не удалось обработать.

Например, событие неизвестной версии отправляется в notifications.dead-letter. Вместе с исходными данными полезно сохранить:

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

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

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

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

Как безопасно перенести сообщение

Пересылка в retry или DLQ тоже состоит из двух действий: записать сообщение в новый топик и сохранить offset исходного.

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

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

Если оба действия выполняются внутри Kafka, их можно объединить в транзакцию: публикацию в целевой топик и commit исходного offset. Это требует корректной транзакционной обработки; одного acks=all для такого переноса недостаточно.

Pause и Resume: переждать общий сбой

Если база полностью недоступна, отправлять каждую новую запись в retry может быть бессмысленно: все они получат одинаковую ошибку.

Consumer может временно приостановить получение новых данных из назначенных партиций через pause(), а после восстановления вызвать resume().

При этом нужно продолжать регулярно вызывать poll(). Пауза не отменяет правила участия в группе. Также нужно сохранить необработанные записи, которые уже были получены до паузы, и не продвинуть offset за них.

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

Избегайте слишком долгой обработки между Poll

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

Для Java consumer этот интервал ограничивает max.poll.interval.ms. Он отличается от механизма heartbeat: наличие сетевых сигналов о жизнеспособности не позволяет бесконечно задерживать poll().

Например, consumer получил 500 записей, а обработка каждой занимает одну секунду. Последовательная обработка займёт больше восьми минут. При лимите между вызовами poll в пять минут это проблема.

После потери назначения другой участник группы может начать работу с последнего сохранённого offset. Первый consumer тем временем ещё выполняет операции. Так возникают повторы, а его попытка commit может завершиться ошибкой.

Если новый consumer обрабатывает ту же пачку так же долго, ситуация повторяется. Частые перераспределения партиций называют `Rebalance Storm).

Что изменить в первую очередь

Сначала нужно измерить время обработки и уменьшить max.poll.records — максимальное число записей, возвращаемых одним poll. Например, вместо 500 проверить 20–50 записей. Конкретное число зависит от времени самой медленной операции.

Это ограничивает число записей, переданных приложению за один poll, а не размер producer-батча или весь объём сетевого чтения.

Затем можно подобрать max.poll.interval.ms с запасом относительно нормального времени обработки. Увеличение лимита не исправляет зависший API или бесконечные retries: у внешних операций тоже должны быть таймауты.

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

Для классического протокола группы один из вариантов — CooperativeStickyAssignor. Но нельзя считать его активным только потому, что он присутствует в списке настроек по умолчанию. Нужно проверить фактически используемые протокол, стратегию назначения и версии клиентов. Настройки consumer.

Когда нужен Worker Pool

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

Представим, что записи 100 и 101 обрабатываются одновременно. Запись 101 уже завершилась, а 100 ещё выполняется. Сохранить позицию 102 пока нельзя: после сбоя запись 100 будет пропущена.

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

Кроме того, нужны ограниченная очередь задач, pause при переполнении и правила завершения работы при передаче партиций другому consumer. Java KafkaConsumer не предназначен для произвольного одновременного использования из нескольких потоков. Многопоточная обработка consumer.

Следите за Consumer Lag

Мы выбрали, как обрабатывать события и переживать ошибки. Теперь нужно видеть, успевает ли система выполнять эту работу.

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

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

Само значение нужно связывать со скоростью обработки. Отставание в 100 000 записей при скорости 50 000 в секунду и при скорости 100 в секунду означает совершенно разную задержку.

Что проверять при росте Lag

Смотрите на каждую партицию, а не только на сумму по группе:

НаблюдениеЧто проверить
Отстают почти все партицииОбщую скорость обработки, базу и внешние API
Отстаёт одна партицияГорячий ключ, тяжёлые сообщения или сбой конкретного consumer
Lag скачет вместе с rebalanceПерезапуски, время между poll и стабильность группы
Lag растёт, consumers отсутствуютЗапуск сервисов, ошибки подключения и назначение партиций

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

Помимо lag полезно измерять бизнес-задержку. Например, сколько времени проходит между оплатой заказа и началом сборки. Поток может формально читаться быстро, а события при этом ждать во внутренней очереди приложения или retry-топике.

Не отправляйте огромные сообщения без необходимости

Большое сообщение нагружает не только сеть между producer и лидером. Оно копируется на followers, хранится во всех репликах и передаётся consumers.

Например, PDF-отчёт размером 100 MB создаёт намного больше нагрузки, чем небольшое событие о его готовности. Кроме того, потребуется согласовать лимиты размера у producer, брокеров и consumers.

Для файлов часто удобнее паттерн Claim Check: файл сохраняется в объектном хранилище, а в Kafka отправляются его идентификатор и ссылка.

{
  "eventType": "ReportReady",
  "eventId": "report-ready-48291",
  "fileId": "report-48291",
  "location": "s3://reports/report-48291.pdf",
  "sizeBytes": 104857600
}

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

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

Базовый production-набор настроек

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

ОбластьНастройкаЧто обеспечивает
Topicreplication.factor=3Три копии каждой партиции
Topic/Brokermin.insync.replicas=2Минимум две реплики в ISR для успешной записи с acks=all
Produceracks=allПодтверждение после репликации в текущем ISR
Producerenable.idempotence=trueЗащиту от дублей при автоматических retries producer
Producerdelivery.timeout.msОграничение времени ожидания результата отправки
Producercompression.type=lz4 или zstdВозможность снизить объём данных ценой работы CPU
Producerlinger.ms и batch.sizeУправление накоплением сообщений в батчи
Consumerenable.auto.commit=falseПередачу контроля над commit приложению
Consumermax.poll.recordsОграничение числа записей за один poll
Consumermax.poll.interval.msДопустимый интервал между вызовами poll
Topicretention.ms и retention.bytesОграничения доступной истории
Topiccleanup.policyВыбор удаления старых сегментов и/или компакции

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

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

Наблюдаемость

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

СигналПочему важен
Consumer Lag и бизнес-задержкаОбработка отстаёт от потребностей приложения
Under-replicated partitionsВ ISR меньше реплик, чем предусмотрено: запас отказоустойчивости снизился
Under-min-ISR partitionsРазмер ISR ниже минимума для записи с acks=all
Offline partitionsУ партиций нет доступного лидера
Частые сокращения и расширения ISRРеплики регулярно отстают или теряют связь
Produce/fetch latency и ошибкиЗапись или чтение замедляются либо завершаются неуспешно
Частые rebalanceНазначения нестабильны, возможна повторная обработка
Размер и возраст outbox/retryСобытия слишком долго ждут публикации или обработки
Новые сообщения в DLQЕсть операции, которые требуют разбора
Свободное место на дискахКластер приближается к пределу хранения

Названия и доступность метрик зависят от версии Kafka и используемого протокола. Метрики Kafka.

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

У каждого такого сигнала должны быть ответственный и понятный первый шаг проверки. Иначе метрика останется красивым графиком без действий.

Пример production-схемы

Соберём решения в процесс обработки заказа:

Order Service
      |
      | одна транзакция: заказ + событие
      v
PostgreSQL / Outbox
      |
      | publisher, Key=orderId
      v
commerce.orders.events
      |
      +--------------------------+
      |                          |
      v                          v
Inventory Group          Notification Group
      |                          |
      v                          v
Резервирование товара     Отправка письма
                                 |
                                 | временная ошибка
                                 v
                         notifications.retry
                                 |
                                 | попытки исчерпаны
                                 v
                         notifications.dead-letter

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

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

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

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

Главные ошибки при внедрении Kafka

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

ОшибкаПоследствие
Считать acks подтверждением бизнес-обработкиСобытие сохранено, но действие не выполнено
Сохранять offset раньше результатаПосле перезапуска необработанная запись может быть пропущена
Не учитывать повторыПовторно начисляются бонусы, резервируется товар или отправляются письма
Генерировать новый eventId при каждой попыткеПолучатель не распознаёт повтор того же события
Переносить запись в retry, забыв о порядкеПозднее событие обрабатывается раньше того, от которого зависит
Складывать ошибки в DLQ без разбораНевыполненные операции накапливаются незаметно
Игнорировать outbox, lag и состояние ISRТеряются время восстановления и запас отказоустойчивости

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

Финал серии

Мы прошли путь от устройства журнала до проектирования потока и восстановления после сбоев. Теперь за словами «Kafka хранит и доставляет сообщения» можно видеть конкретные этапы: публикацию, репликацию, чтение, бизнес-обработку и сохранение позиции.

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

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