Kafka: Retry, DLQ и эксплуатация Consumer — Часть 9

Содержание
Коротко: Повторы должны иметь понятный предел и учитывать порядок событий. Retry и DLQ требуют надёжного переноса сообщений, а lag, состояние ISR и бизнес-задержка показывают, успевает ли система восстановиться.

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

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

В этой части я предлагаю посмотреть на поток глазами команды, которой нужно вернуть его в работу. Сначала разберём retry, DLQ и паузу чтения, затем — ограничения poll(), отставание и признаки проблем в работающей системе. Нас будет интересовать не только «как повторить», но и «как понять, что повторы вообще помогают».

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

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

Я бы начинал разбор не с количества повторов, а с причины ошибки. Для этого полезно различать:

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

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

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

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

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

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

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

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

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

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

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

commerce.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. Вместе с исходными данными полезно сохранить:

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

Следите за Consumer Lag

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

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

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

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

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

При росте lag я рекомендую сначала посмотреть на отдельные Partition, а не только на сумму по группе. Это помогает отличить общую нехватку производительности от проблемы одного обработчика или ключа:

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

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

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

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

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

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

Для файлов часто удобнее паттерн 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Три копии каждой Partition
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.

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

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

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

СигналПочему важен
Consumer Lag и бизнес-задержкаОбработка отстаёт от потребностей приложения
Under-replicated partitionsВ ISR меньше реплик, чем предусмотрено: запас отказоустойчивости снизился
Under-min-ISR partitionsРазмер ISR ниже минимума для записи с acks=all
Offline partitionsУ Partition нет доступного лидера
Частые сокращения и расширения 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 Group и читают поток независимо. Событие сохраняется в outbox вместе с заказом, а publisher передаёт его в Kafka с подтверждением записи.

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

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

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

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

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

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

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

Финал серии

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

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

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