📨 Розділ 15 · Питання #25

Що таке 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

  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 ніяк не впливає на порядок повідомлень» — Порушує суворий порядок за ключем партиціонування при подальшому реплеї.

Пов’язані теми