📨 Раздел 15 · Вопрос #25

Что такое DLQ (Dead Letter Queue)

В традиционных брокерах (RabbitMQ, AWS SQS) механизм DLQ встроен в сам сервер очередей. В Apache Kafka брокер не имеет встроенной концепции DLQ — партиция является неизменяемым...


🟢 Junior Level

Что такое Dead Letter Queue простыми словами

Dead Letter Queue (DLQ, «очередь недоставленных / мертвых писем») — это специализированный топик в Apache Kafka (часто называемый DLT — Dead Letter Topic), в который перенаправляются сообщения, которые консьюмер не смог успешно обработать после исчерпания всех попыток повтора (retries).

В традиционных брокерах (RabbitMQ, AWS SQS) механизм DLQ встроен в сам сервер очередей. В Apache Kafka брокер не имеет встроенной концепции DLQ — партиция является неизменяемым логом. Поэтому DLQ в Kafka — это архитектурный паттерн клиентского уровня, реализуемый либо вручную в коде консьюмера, либо через фреймворки (Spring Kafka, Kafka Connect).

Основной поток обработки:
[Топик: orders] ──> Consumer ──> Попытка 1 (Сбой)
                                  Попытка 2 (Сбой)
                                  Попытка 3 (Сбой)
                                        │
                                        ▼ (Исчерпаны ретраи)
[Топик: orders.DLT] <── Продюсер консьюмера отправляет сообщение в DLQ
[Топик: orders]     ──> Консьюмер коммитит исходный offset и переходит к следующему сообщению!

Зачем нужен DLQ

  1. Защита от «ядовитых пилюль» (Poison Pills): Если в топик попало сообщение с невалидным JSON или поврежденной схемой, без DLQ консьюмер будет падать на нем вечно, блокируя обработку всех последующих валидных сообщений.
  2. Изоляция сбоев без остановки конвейера: DLQ позволяет закоммитить офсет сбойного сообщения в основном топике, переместив «брак» в отдельное место, чтобы консьюмер продолжил обрабатывать очередь.
  3. Сохранение данных для расследования: Ни одно сообщение не теряется. Инженеры могут исследовать причину сбоя, исправить баг в коде и заново отправить сообщения из DLQ в рабочий топик (Replay).

Базовый пример настройки в Spring Kafka

@Configuration
public class KafkaListenerConfig {

    @Bean
    public DefaultErrorHandler errorHandler(KafkaOperations<Object, Object> template) {
        // Отправляет сбойное сообщение в топик <originalTopic>.DLT с обогащением заголовков
        DeadLetterPublishingRecoverer recoverer = new DeadLetterPublishingRecoverer(template);
        
        // 3 попытки с интервалом в 1 секунду перед отправкой в DLT
        FixedBackOff backOff = new FixedBackOff(1000L, 3);
        
        return new DefaultErrorHandler(recoverer, backOff);
    }
}

🟡 Middle Level

Обогащение метаданными (Header Enrichment)

Просто переслать тело (payload) сбойного сообщения в DLQ — грубая ошибка проектирования. Инженер, разбирающий инцидент, не сможет понять, почему сообщение упало и откуда оно взялось.

При публикации в DLQ фреймворк или прикладной код обязан обогатить запись служебными заголовками (Headers):

public ProducerRecord<String, String> enrichForDlq(
        ConsumerRecord<String, String> originalRecord, 
        Exception exception) {
    
    ProducerRecord<String, String> dlqRecord = new ProducerRecord<>(
        originalRecord.topic() + ".DLT",
        originalRecord.key(),
        originalRecord.value()
    );

    // Добавление критически важных диагностических заголовков
    Headers headers = dlqRecord.headers();
    headers.add("x-original-topic", originalRecord.topic().getBytes(StandardCharsets.UTF_8));
    headers.add("x-original-partition", ByteBuffer.allocate(4).putInt(originalRecord.partition()).array());
    headers.add("x-original-offset", ByteBuffer.allocate(8).putLong(originalRecord.offset()).array());
    headers.add("x-original-timestamp", ByteBuffer.allocate(8).putLong(originalRecord.timestamp()).array());
    headers.add("x-exception-fqcn", exception.getClass().getName().getBytes(StandardCharsets.UTF_8));
    headers.add("x-exception-message", String.valueOf(exception.getMessage()).getBytes(StandardCharsets.UTF_8));
    headers.add("x-failed-at-host", getHostName().getBytes(StandardCharsets.UTF_8));

    return dlqRecord;
}

В Spring Kafka стандартный DeadLetterPublishingRecoverer делает это автоматически, добавляя заголовки:

  • KafkaHeaders.DLT_ORIGINAL_TOPIC
  • KafkaHeaders.DLT_ORIGINAL_PARTITION
  • KafkaHeaders.DLT_ORIGINAL_OFFSET
  • KafkaHeaders.DLT_EXCEPTION_FQCN
  • KafkaHeaders.DLT_EXCEPTION_MESSAGE
  • KafkaHeaders.DLT_EXCEPTION_STACKTRACE

Архитектурные паттерны организации DLQ

Паттерн Схема именования Плюсы Минусы
Единый топик на приложение (Shared DLQ) app-common.DLT Простота создания и единая точка мониторинга Сложно десериализовать разные схемы данных; разные SLA
Топик на каждый бизнес-топик (Per-Topic DLT) orders.DLT, payments.DLT Рекомендуемый стандарт. Изоляция форматов и прав доступа Больше топиков и партиций в кластере
По типам ошибок (Per-Error DLT) orders-validation.DLT, orders-timeout.DLT Автоматическая маршрутизация: битые схемы — разработчикам, таймауты — в авто-реплей Избыточная сложность топологии

Настройки хранения топика DLQ

Топики DLQ требуют специальной конфигурации на брокере:

  1. Retention Period (retention.ms): Нельзя оставлять бесконечным. Обычно выставляется от 7 до 30 дней. Это дает инженерам время на расследование, исключая переполнение дисков брокера.
  2. Cleanup Policy: Строго delete. Компактизация (compact) для DLQ недопустима, так как она сотрет историю повторяющихся ошибок по одному и тому же ключу.
  3. Replication Factor: Должен быть равен основному топику (минимум 3 в проде) во избежание потери инцидентных данных при падении брокера.

🔴 Senior Level

Проблема нарушения причинно-следственной связи (Causality & Ordering Breakage)

Самый опасный скрытый побочный эффект вытеснения сообщения в DLQ — нарушение строгого порядка бизнес-событий (Ordering Guarantee):

Партиция orders (Ключ = OrderID 777):
Offset 10: Event "CREATED"  ──> Обработано успешно (Заказ создан в БД)
Offset 11: Event "PAYMENT"  ──> Сбой (Timeout платежного шлюза) ──> Отправлен в DLQ!
Offset 12: Event "CANCELLED"──> Обработано успешно (Заказ отменен в БД)

Позже инженер делает Replay сообщения Offset 11 из DLQ:
Event "PAYMENT" исполняется после Event "CANCELLED"!
Итог: С пользователя списаны деньги за уже отмененный заказ (нарушение инварианта системы).

Архитектурные решения проблемы:

  1. Circuit Breaker / Stop Consumption: Для критичных бизнес-доменов (финансы, складской учет) вытеснение в DLQ запрещено. При сбое консьюмер останавливает чтение партиции (consumer.pause()), алерт вызывает дежурного инженера, порядок не нарушается.
  2. Паттерн Saga / Idempotent State Machine: Сервис обработки проверяет статус агрегата. Если заказ находится в статусе CANCELLED, запоздалое событие PAYMENT из DLQ отвергается логикой бизнес-валидатора.

Стратегии повторной обработки (DLQ Replay Pipeline)

Отправка сообщения в DLQ — это только половина решения. Вторая половина — как вернуть эти сообщения в строй:

                                  DLQ REPLAY ARCHITECTURE
                                  
   [orders.DLT] 
        │
        ▼
   DLQ Consumer Service (или UI Tool: AKHQ, Provectus UI, Kpow)
        │
        ├── 1. Автоматический реплей (для устраненных транзитных сбоев):
        │      Считывает заголовки, проверяет здоровье внешнего сервиса,
        │      очищает служебные DLT-заголовки и шлет в исходный топик "orders".
        │
        ├── 2. Ручной патч (Manual Fix):
        │      Аналитик/разработчик правит поврежденное JSON-поле в UI
        │      и отправляет исправленную версию в топик.
        │
        └── 3. DLQ-for-DLQ (Parking Lot / Quarantine):
               Если сообщение при реплее снова падает — оно отправляется
               в финальный карантин orders-parking-lot для ручного списания.

Транзакционная безопасность при публикации в DLQ

Чтобы исключить потерю сообщения или его дублирование в DLQ при аварии консьюмера:

  • Консьюмер должен объединять чтение из входящего топика, публикацию в DLT и коммит смещений в единую Kafka-транзакцию:
// Псевдокод транзакционного Consumer-Producer цикла
producer.beginTransaction();
try {
    process(record);
} catch (Exception fatalError) {
    ProducerRecord<String, String> dlqRecord = enrich(record, fatalError);
    producer.send(dlqRecord);
}
// Смещение коммитится строго внутри транзакции
producer.sendOffsetsToTransaction(
    Map.of(new TopicPartition(record.topic(), record.partition()), new OffsetAndMetadata(record.offset() + 1)),
    consumerGroupMetadata
);
producer.commitTransaction();

4 Tricky Questions

1. Поддерживает ли Apache Kafka брокер встроенный DLQ на уровне самого сервера (как RabbitMQ или AWS SQS)?

Ответ: Нет, в ядре Apache Kafka нет концепции DLQ. В RabbitMQ брокер сам проверяет счетчик доставки сообщения и при basic.nack(requeue=false) маршрутизирует сообщение через dead-letter-exchange в другую очередь без участия клиента. В Kafka брокер — это пассивный append-only лог, не знающий о статусе бизнес-обработки записей консьюмерами. DLQ в Kafka реализуется исключительно на стороне клиента (Consumer читает сбойное сообщение, выступает продюсером для топика .DLT, и только после успешной отправки сдвигает свой commit offset).

2. Что произойдет, если топик DLQ недоступен или отклонил сообщение (например, превышен message.max.bytes)?

Ответ:

  1. Вызов template.send(dlqRecord).get() завершится исключением (TimeoutException, RecordTooLargeException или KafkaException).
  2. В правильно спроектированном обработчике ошибок (DefaultErrorHandler) это приведет к пробросу исключения в контейнер слушателя.
  3. Офсет сбойного сообщения в основном топике не будет закоммичен.
  4. Консьюмер остановится либо повторит попытку чтения того же смещения на следующей итерации poll().
  5. Антипаттерн: Если перехватить исключение отправки в DLQ в пустой блок catch и закоммитить офсет — произойдет фатальная потеря данных («Silent Data Loss»).

3. Как Spring Kafka DeadLetterPublishingRecoverer выбирает целевую партицию в DLQ топике, и к какой проблеме это может привести?

Ответ: По умолчанию DeadLetterPublishingRecoverer пытается отправить сообщение в ту же самую по номеру партицию, из которой оно было вычитано в оригинальном топике: \(Partition_{DLQ} = Partition_{Original}\)

  • Ловушка конфигурации: Если топик orders имеет 16 партиций, а созданный топик orders.DLT имеет всего 3 партиции (настройка по умолчанию для автосоздаваемых топиков), отправка сообщений из партиций 3–15 немедленно упадет с ошибкой KafkaException: Partition 7 is not valid for topic orders.DLT!
  • Решение: Либо создавать DLQ-топик с тем же числом партиций, либо переопределить резолвер назначения, направляя сообщения по стандартному партиционированию ключа или в партицию 0: new DeadLetterPublishingRecoverer(template, (rec, ex) -> new TopicPartition(rec.topic() + ".DLT", -1)).

4. В чем заключается феномен «DLQ Storm» и как предотвратить исчерпание дискового пространства кластера?

Ответ: DLQ Storm (шторм недоставленных писем) возникает, когда внешний сервис (например, главная реляционная БД или биллинг) падает полностью, либо продюсер начинает массово слать невалидные сообщения:

  1. Миллионы входящих сообщений подряд исчерпывают лимит ретраев и лавинообразно направляются в топик DLQ.
  2. Топик DLQ начинает расти со скоростью основного потока (сотни тысяч RPS).
  3. Дисковые массивы брокеров переполняются, брокеры переходят в read-only режим или падают по аварийному останову диска.
  4. Защитные меры:
    • Внедрение паттерна Circuit Breaker (Resilience4j): если процент ошибок превышает 50%, консьюмер не шлет сообщения в DLQ, а временно приостанавливает вызовы poll() основного топика.
    • Жесткий мониторинг скорости прироста DLQ (rate(kafka_topic_messages_total{topic=~".*DLT"}[1m]) > threshold).
    • Настройка строгой политики retention на DLQ-топиках.

🎯 Шпаргалка для интервью

30-секундный ответ (Elevator Pitch)

«Dead Letter Queue (DLQ / DLT) в Kafka — это архитектурный паттерн клиентского уровня для изоляции сообщений, которые невозможно обработать после заданного числа попыток. Поскольку лог Kafka неизменяем, консьюмер сам публикует сбойную запись в специальный топик (например, <topic>.DLT), обогащая ее заголовками с причиной ошибки, стектрейсом и исходным офсетом, после чего коммитит смещение в исходном топике. Это спасает консьюмер от зависания на poison pill сообщениях. Ключевые риски DLQ — нарушение строгой причинно-следственной последовательности событий по одному ключу и риск шторма переполнения дисков (DLQ Storm) при падении внешних зависимостей».

Ключевые метрики для Production-мониторинга

  • kafka_topic_messages_in_total{topic=~".*DLT"}: скорость поступления сообщений в топики DLQ. Любой всплеск — повод для немедленного PagerDuty/Telegram алерта.
  • kafka_consumergroup_lag{topic=~".*DLT"}: лаг сервиса-обработчика DLQ.
  • disk_used_percent на брокерах, хранящих DLQ-топики.

Красные флаги (Чего говорить НЕЛЬЗЯ)

  • ❌ «Брокер Kafka сам автоматически перекладывает битые сообщения в DLQ, как в RabbitMQ» — В ядре Kafka нет механизма DLQ; это логика прикладного консьюмера или фреймворка.
  • ❌ «В DLQ достаточно отправить только payload сообщения» — Без заголовков с исходным топиком, смещением, причиной и стектрейсом разобрать инцидент невозможно.
  • ❌ «DLQ можно включить и забыть, он никому не мешает» — Без мониторинга и регламента очистки/реплея DLQ приведет к переполнению диска или скрытой потере бизнес-данных.
  • ❌ «Вытеснение в DLQ никак не влияет на порядок сообщений» — Нарушает строгий порядок по ключу партиционирования при последующем реплее.

Связанные темы