Що таке DLQ (Dead Letter Queue)
У традиційних брокерах (RabbitMQ, AWS SQS) механізм DLQ вбудований у сам сервер черг. В Apache Kafka брокер не має вбудованої концепції DLQ — партиція є незмінним логом. Тому DL...
🟢 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 ніяк не впливає на порядок повідомлень» — Порушує суворий порядок за ключем партиціонування при подальшому реплеї.