Что такое 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
- Защита от «ядовитых пилюль» (Poison Pills): Если в топик попало сообщение с невалидным JSON или поврежденной схемой, без DLQ консьюмер будет падать на нем вечно, блокируя обработку всех последующих валидных сообщений.
- Изоляция сбоев без остановки конвейера: DLQ позволяет закоммитить офсет сбойного сообщения в основном топике, переместив «брак» в отдельное место, чтобы консьюмер продолжил обрабатывать очередь.
- Сохранение данных для расследования: Ни одно сообщение не теряется. Инженеры могут исследовать причину сбоя, исправить баг в коде и заново отправить сообщения из 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_TOPICKafkaHeaders.DLT_ORIGINAL_PARTITIONKafkaHeaders.DLT_ORIGINAL_OFFSETKafkaHeaders.DLT_EXCEPTION_FQCNKafkaHeaders.DLT_EXCEPTION_MESSAGEKafkaHeaders.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 требуют специальной конфигурации на брокере:
- Retention Period (
retention.ms): Нельзя оставлять бесконечным. Обычно выставляется от 7 до 30 дней. Это дает инженерам время на расследование, исключая переполнение дисков брокера. - Cleanup Policy: Строго
delete. Компактизация (compact) для DLQ недопустима, так как она сотрет историю повторяющихся ошибок по одному и тому же ключу. - 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"!
Итог: С пользователя списаны деньги за уже отмененный заказ (нарушение инварианта системы).
Архитектурные решения проблемы:
- Circuit Breaker / Stop Consumption: Для критичных бизнес-доменов (финансы, складской учет) вытеснение в DLQ запрещено. При сбое консьюмер останавливает чтение партиции (
consumer.pause()), алерт вызывает дежурного инженера, порядок не нарушается. - Паттерн 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)?
Ответ:
- Вызов
template.send(dlqRecord).get()завершится исключением (TimeoutException,RecordTooLargeExceptionилиKafkaException). - В правильно спроектированном обработчике ошибок (
DefaultErrorHandler) это приведет к пробросу исключения в контейнер слушателя. - Офсет сбойного сообщения в основном топике не будет закоммичен.
- Консьюмер остановится либо повторит попытку чтения того же смещения на следующей итерации
poll(). - Антипаттерн: Если перехватить исключение отправки в 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 (шторм недоставленных писем) возникает, когда внешний сервис (например, главная реляционная БД или биллинг) падает полностью, либо продюсер начинает массово слать невалидные сообщения:
- Миллионы входящих сообщений подряд исчерпывают лимит ретраев и лавинообразно направляются в топик DLQ.
- Топик DLQ начинает расти со скоростью основного потока (сотни тысяч RPS).
- Дисковые массивы брокеров переполняются, брокеры переходят в read-only режим или падают по аварийному останову диска.
- Защитные меры:
- Внедрение паттерна Circuit Breaker (Resilience4j): если процент ошибок превышает 50%, консьюмер не шлет сообщения в DLQ, а временно приостанавливает вызовы
poll()основного топика. - Жесткий мониторинг скорости прироста DLQ (
rate(kafka_topic_messages_total{topic=~".*DLT"}[1m]) > threshold). - Настройка строгой политики retention на DLQ-топиках.
- Внедрение паттерна Circuit Breaker (Resilience4j): если процент ошибок превышает 50%, консьюмер не шлет сообщения в 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 никак не влияет на порядок сообщений» — Нарушает строгий порядок по ключу партиционирования при последующем реплее.