Retry-топик и DLQ в Kafka

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

Kafka не умеет «пропустить» сообщение. Консьюмер читает партицию строго последовательно, и единственный способ перешагнуть запись — закоммитить оффсет, то есть признать её обработанной. Поэтому запись, которую обработчик не может переварить, останавливает всю партицию за собой. Это называется head-of-line blocking, и в проде выглядит как «интеграция встала», хотя кластер полностью здоров.

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

После исчерпания попыток запись уходит в DLQ (dead letter queue) — топик недоставленных сообщений. В DLQ важно класть не только тело, но и исходные заголовки, причину ошибки, стектрейс и число попыток: без этого разбор превращается в гадание.

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

Альтернатива для случаев, где порядок не важен: share groups (Kafka 4.2+). В них у каждой записи свой счётчик попыток доставки, и после исчерпания лимита запись уходит в archived, не блокируя остальные.

Разобрать это в симуляторе Интерактивный сценарий: настраиваешь кластер, ломаешь его и смотришь, что происходит

Частые вопросы

Сколько попыток делать?

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

Нужен ли отдельный retry-топик на каждую задержку?

Частая схема — retry-топики по уровням задержки (5s, 1m, 10m). Проще в реализации, чем отложенная обработка внутри одного топика.

Можно ли обойтись без DLQ?

Можно, если бизнес готов терять такие сообщения и это зафиксировано. Молча отбрасывать без решения — нельзя.

Дальше по теме

Что такое партиция в Kafka Consumer lag: что это и как с ним работать acks в Kafka: 0, 1 и all ISR и min.insync.replicas Ребалансировка consumer group