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

Як обробляти помилки при читанні повідомлень

Під час читання повідомлень із Kafka збої неминучі: відмова зовнішньої бази даних, мережевий таймаут HTTP-клієнта, невалідний формат JSON або порушення бізнес-правил.


🟢 Junior Level

Суть проблеми обробки помилок консьюмера

Під час читання повідомлень із Kafka збої неминучі: відмова зовнішньої бази даних, мережевий таймаут HTTP-клієнта, невалідний формат JSON або порушення бізнес-правил.

На відміну від класичних черг (RabbitMQ, SQS), де окреме повідомлення можна просто повернути назад у початок черги (nack / reject), у Kafka партиція є незмінним упорядкованим логом (Append-only Log). Неможливо «вилучити» збійне повідомлення із середини партиції і пропустити його, не зміщуючи загальний вказівник зсуву (offset).

Лог партиції:
[Offset 100: OK] -> [Offset 101: ERROR (Poison Pill)] -> [Offset 102: OK] -> [Offset 103: OK]
                         │
                         ▼
        Що робити з вказівником commit offset?
        1. Закомітити 101 -> ВТРАТА ДАНИХ (повідомлення не оброблено).
        2. Не комітити і впасти -> НЕСКІНЧЕННИЙ ЦИКЛ (при рестарті знову впадемо на 101).

Основні сценарії та базовий підхід

  1. Транзитні помилки (Transient Errors): Короткочасні мережеві збої або блокування БД. Вирішуються кількома повторними спробами (Retry).
  2. Фатальні помилки («Отруйні пігулки», Poison Pills): Пошкоджені байти, несумісні схеми Protobuf/Avro, ділення на нуль. Повтори безглузді, повідомлення відправляється в Dead Letter Queue (DLQ) для подальшого розбору.

Базовий код безпечної обробки на чистому Java Kafka Consumer

Properties props = new Properties();
props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");
props.put(ConsumerConfig.GROUP_ID_CONFIG, "order-processors");
props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class.getName());
props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class.getName());
// 1. Вимикаємо автоматичний коміт для повного контролю помилок
props.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, "false");

KafkaConsumer<String, String> consumer = new KafkaConsumer<>(props);
consumer.subscribe(List.of("orders"));

try {
    while (running) {
        ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(500));
        
        for (ConsumerRecord<String, String> record : records) {
            boolean success = false;
            int attempts = 0;
            
            while (!success && attempts < 3) {
                try {
                    attempts++;
                    processOrder(record.value());
                    success = true;
                } catch (TransientDatabaseException ex) {
                    log.warn("Тимчасовий збій БД на зсуві {}, спроба {}", record.offset(), attempts);
                    // Коротка пауза перед повтором (увага: не можна перевищувати max.poll.interval.ms!)
                    Thread.sleep(1000);
                } catch (Exception fatalEx) {
                    log.error("Фатальна помилка обробки. Відправка в DLQ: offset={}", record.offset(), fatalEx);
                    sendToDlq("orders-dlq", record, fatalEx);
                    success = true; // Вважаємо обробленим через DLQ, щоб не блокувати партицію
                    break;
                }
            }
            
            if (!success) {
                // Вичерпано ліміт ретраїв: скидаємо в DLQ
                sendToDlq("orders-dlq", record, new RuntimeException("Retries exhausted"));
            }
        }
        
        // 2. Комітимо зсуви лише після того, як усі записи батча успішно оброблені або відправлені в DLQ
        consumer.commitSync();
    }
} catch (WakeupException e) {
    log.info("Завершення роботи консьюмера...");
} finally {
    consumer.close();
}

🟡 Middle Level

Класифікація помилок і стратегії їх обробки

Тип помилки Приклади винятків Стратегія Вплив на порядок
Десеріалізація (Corrupted Payload) SerializationException, пошкоджений JSON/Avro ErrorHandlingDeserializer + DLQ Порядок зберігається, бітий оффсет пропускається
Транзитний збій інфраструктури SocketTimeoutException, LockAcquisitionException Блокуючий Retry з паузою Порядок суворо зберігається, партиція чекає відновлення
Бізнес-помилка валідації Від’ємна сума замовлення, невідомий статус DLQ (без ретраїв) Повідомлення ізолюється, потік продовжується
Фатальний збій сервісу OutOfMemoryError, втрата прав доступу Аварійна зупинка (Stop Container) Консьюмер падає, алерт черговому інженеру

Пастка max.poll.interval.ms при блокуючих ретраях

Якщо всередині циклу консьюмера виконувати тривалі паузи Thread.sleep() або нескінченні ретраї:

  1. Консьюмер не встигає викликати наступний consumer.poll() за час max.poll.interval.ms (за замовчуванням 300 000 мс = 5 хвилин).
  2. Group Coordinator на брокері вважає консьюмер «мертвим» (завислим).
  3. Ініціюється Rebalance: партиція відбирається і передається іншому консьюмеру групи.
  4. Новий консьюмер починає читати з останнього закоміченого оффсету, натикається на те саме збійне повідомлення, знову зависає в ретраях і знову вилітає за таймаутом.
  5. Результат: Каскадний шторм ребалансувань (Rebalance Storm), повна зупинка всієї групи.

Архітектура Spring Kafka: DefaultErrorHandler

У сучасних Java-застосунках на Spring Boot для обробки помилок використовується DefaultErrorHandler (прийшов на зміну застарілому SeekToCurrentErrorHandler у Spring Kafka 2.8+):

@Configuration
public class KafkaConsumerConfig {

    @Bean
    public DefaultErrorHandler errorHandler(KafkaTemplate<Object, Object> template) {
        // 1. Налаштування DLQ-відновлювача: публікує збійне повідомлення в топік <originalTopic>.DLT
        DeadLetterPublishingRecoverer recoverer = new DeadLetterPublishingRecoverer(template,
            (record, exception) -> new TopicPartition(record.topic() + ".DLT", record.partition()));

        // 2. Експоненційний BackOff: 1s, 2s, 4s, максимум 3 спроби
        ExponentialBackOffWithMaxRetries backOff = new ExponentialBackOffWithMaxRetries(3);
        backOff.setInitialInterval(1000L);
        backOff.setMultiplier(2.0);
        backOff.setMaxInterval(10000L);

        DefaultErrorHandler errorHandler = new DefaultErrorHandler(recoverer, backOff);

        // 3. Неретраябельні винятки відправляються в DLQ миттєво без пауз
        errorHandler.addNotRetryableExceptions(
            IllegalArgumentException.class,
            JsonParseException.class,
            NoSuchElementException.class
        );

        return errorHandler;
    }
}

Як DefaultErrorHandler вирішує проблему таймаутів: Він не присипляє потік через Thread.sleep(). Натомість він перехоплює виняток, здійснює consumer.pause(partitions), робить seek() назад на проблемний оффсет і виконує короткі швидкі poll(0), надсилаючи брокеру heartbeat-пакети. Після закінчення таймера паузи партиція відновлюється через consumer.resume().


🔴 Senior Level

Неблокуючі ретраї: Паттерн Retry Topics (Uber / Spring @RetryableTopic)

У високонавантажених системах (Highload) блокувати партицію на 5–10 хвилин заради очікування відновлення зовнішнього CRM/Billing сервісу неприпустимо — відставання (lag) виросте на мільйони повідомлень.

Для вирішення проблеми застосовується неблокуючий конвеєр черг затримок (Non-Blocking Delay Topics):

[Основний топік: orders] ──> Consumer 1 (Збій виклику зовнішнього API)
                                  │
                                  ▼
           [Топік orders-retry-10s] ──> Consumer 2 (Читає із затримкою 10 сек)
                                             │
                                             ▼ (Знову збій)
           [Топік orders-retry-60s] ──> Consumer 3 (Читає із затримкою 60 сек)
                                             │
                                             ▼ (Усі спроби вичерпано)
           [Топік orders-dlq]       ──> Ручний розбір / Дашборд техпідтримки

Реалізація через Spring @RetryableTopic

@Component
public class OrderEventsListener {

    @RetryableTopic(
        attempts = "4",
        backoff = @Backoff(delay = 10000, multiplier = 2.0, maxDelay = 60000),
        autoCreateTopics = "false",
        topicSuffixingStrategy = TopicSuffixingStrategy.SUFFIX_WITH_DELAY_VALUE,
        dltStrategy = DltStrategy.FAIL_ON_ERROR,
        include = { RemoteServiceUnavailableException.class, TimeoutException.class }
    )
    @KafkaListener(topics = "orders", groupId = "order-billing-group")
    public void handleOrder(OrderEvent event, @Header(KafkaHeaders.RECEIVED_TOPIC) String topic) {
        log.info("Обробка замовлення {} з топіка {}", event.orderId(), topic);
        billingClient.charge(event);
    }

    @DltHandler
    public void processDlt(OrderEvent event, @Header(KafkaHeaders.ORIGINAL_OFFSET) long offset) {
        log.error("Замовлення {} остаточно провалило обробку, offset={}. Потрібне втручання інженера.",
            event.orderId(), offset);
        alertingService.notifyOnCallSupport(event);
    }
}

Критичний архітектурний компроміс (Trade-off): Перенесення збійних повідомлень в окремі retry-топіки порушує суворий порядок обробки повідомлень (Ordering Guarantee) у межах одного бізнес-ключа (key). Якщо за користувачем надійшла подія «Створити замовлення», а слідом «Скасувати замовлення», то перенесення першої в ретрай-топік може призвести до того, що скасування виконається раніше за створення.

Проблема Poison Pill на етапі десеріалізації (ErrorHandlingDeserializer)

Якщо повідомлення містить пошкоджені байти, стандартний StringDeserializer або KafkaAvroDeserializer викине виняток прямо всередині системного виклику consumer.poll().

  • До бізнес-коду виконання навіть не дійде.
  • Консьюмер перехоплює помилку, відкочується, робить наступний poll() і знову падає на десеріалізації. Партиція наглухо блокується.

Рішення на рівні архітектури Kafka: Використання org.springframework.kafka.support.serializer.ErrorHandlingDeserializer:

key.deserializer=org.apache.kafka.common.serialization.StringDeserializer
# Замість реального десеріалізатора вказуємо ErrorHandlingDeserializer
value.deserializer=org.springframework.kafka.support.serializer.ErrorHandlingDeserializer
# Всередину передаємо цільовий робочий клас десеріалізатора
spring.deserializer.value.delegate.class=org.springframework.kafka.support.serializer.JsonDeserializer

ErrorHandlingDeserializer не викидає виняток при збої парсингу. Натомість він повертає null як значення value(), а в заголовки запису (Headers) упаковує метадані про помилку:

  • DeserializationException.KEY_EXCEPTION_FQCN
  • DeserializationException.VALUE_EXCEPTION_STACKTRACE

Далі DefaultErrorHandler зчитує ці заголовки і одразу скеровує пошкоджене сире повідомлення в DLT, оминаючи бізнес-логіку!


4 Tricky Questions

1. Як реалізувати повторну спробу (Retry) у чистому KafkaConsumer API без використання сторонніх фреймворків і без ризику зловити CommitFailedException / виліт за max.poll.interval.ms?

Відповідь: Для цього використовується механізм паузи партицій:

  1. При отриманні помилки консьюмер зберігає TopicPartition та оффсет впалого запису в тимчасовий стан.
  2. Викликається consumer.pause(Collections.singleton(topicPartition)).
  3. Виконується скидання вказівника: consumer.seek(topicPartition, failedOffset).
  4. В основному циклі застосунку продовжують викликатися consumer.poll(Duration.ofMillis(100)) з нормальною регулярністю (кожні 100–200 мс). Оскільки партиція на паузі, брокер повертає порожні пакети даних, але координатор отримує підтвердження життя потоку (heartbeat), запобігаючи ребалансу.
  5. Фоновий таймер або планувальник відлічує час затримки (backoff). Після закінчення таймера викликається consumer.resume(Collections.singleton(topicPartition)).
  6. Наступний poll() знову прочитає проблемний запис. Порядок повідомлень у партиції дотримується на 100%.

2. Що станеться, якщо в процесі обробки пакета з 500 записів повідомлення № 250 впало з помилкою, а в коді використовується enable.auto.commit = true?

Відповідь: Відбудеться тиха втрата даних (Silent Data Loss):

  1. При enable.auto.commit = true фоновий потік комітить поточний максимальний оффсет пакета кожні auto.commit.interval.ms (за замовчуванням 5 секунд) під час виклику poll().
  2. Якщо виклик poll() повернув записи з 1 по 500, і на записі 250 застосунок викинув виняток, перервав цикл і перейшов до наступного poll(), консьюмер автоматично закомітить оффсет 500.
  3. Повідомлення з 250 по 500 залишаться необробленими, але закоміченими.
  4. При падінні або рестарті процесу читання почнеться з оффсету 501. Записи 250–500 втрачені безповоротно.
  5. Висновок: Для надійних систем enable.auto.commit завжди повинен бути вимкнений (false).

3. Чому відправка в Dead Letter Queue (DLQ) зобов’язана бути синхронною або транзакційною перед комітом зсуву?

Відповідь: Якщо консьюмер відправляє повідомлення в DLQ асинхронно через звичайний dlqProducer.send(record) без очікування результату:

  1. Виклик send() повертає Future і кладе повідомлення в локальний мережевий буфер RecordAccumulator.
  2. Консьюмер тут же виконує consumer.commitSync() основного топіка, вважаючи задачу виконаною.
  3. У цей момент мережа до DLQ-топіка розривається, або DLQ-брокер переповнений (викинуто RecordTooLargeException / TimeoutException).
  4. Буфер продюсера очищається з помилкою у фоновому потоці, а вихідний оффсет вже закомічено в основному топіку.
  5. Повідомлення безповоротно втрачено: його немає ні в основному потоці, ні в DLQ.
  6. Правило: Консьюмер зобов’язаний або дочекатися підтвердження запису в DLQ (dlqProducer.send().get()), або об’єднувати читання, запис у DLQ та коміт оффсету в єдину розподілену транзакцію (sendOffsetsToTransaction).

4. У чому небезпека використання @RetryableTopic для топіків з високим навантаженням (High RPS), у яких ключі партиціонування є критичними для бізнес-логіки?

Відповідь:

  1. Порушення партиціонування за ключами: За замовчуванням допоміжні retry-топіки можуть створюватися з іншою кількістю партицій. Навіть при однаковій кількості партицій алгоритм Sticky/Murmur2 партиціонування гарантує порядок лише всередині однієї черги.
  2. Race Conditions (Гонки станів): Якщо подія $E_1$ (збійна) потрапила в топік із затримкою 10 секунд, а подія $E_2$ за тим самим агрегатом надійшла через секунду й успішно обробилася основним консьюмером, то через 9 секунд консьюмер retry-топіка застосує застарілий стан $E_1$ поверх свіжого $E_2$.
  3. Множення мережевого трафіку (Traffic Amplification): При кожному ретраї повідомлення фізично пересилається продюсером і записується в новий топік на брокерах, множачи дисковий I/O та мережевий egress в $N$ разів (де $N$ — кількість спроб).

🎯 Шпаргалка для інтерв’ю

30-секундна відповідь (Elevator Pitch)

«У Kafka партиція — це незмінний лог, тому не можна пропустити бите повідомлення, просто проігнорувавши його. Стратегія обробки помилок залежить від їх типу: для транзитних збоїв (таймаути БД, мережі) застосовується Retry з експоненційним Backoff; для неретраябельних фатальних помилок (Poison Pills) повідомлення перенаправляється в Dead Letter Queue (DLQ) зі збереженням стек-трейсу та метаданих у заголовках. Для збереження суворого порядку повідомлень ретрай повинен бути блокуючим (з використанням consumer.pause() / seek() для захисту від перевищення max.poll.interval.ms). Якщо суворий порядок за ключем не потрібен, застосовується паттерн Non-blocking Retry Topics (кілька черг затримок). При збоях на етапі десеріалізації обов’язково використовується ErrorHandlingDeserializer».

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

  • records-lag-max: максимальне відставання консьюмера по партиціях топіка.
  • app_kafka_consumer_dlq_total: лічильник повідомлень, відправлених у топіки DLQ (тригер для чергових алертів).
  • app_kafka_consumer_retry_rate: відсоток повідомлень, що потребують повторної обробки.
  • last-heartbeat-seconds-ago: час з моменту останнього heartbeat. Різке зростання свідчить про блокування консьюмера і загрозу ребалансу.

Червоні прапорці (Чого говорити НЕ МОЖНА)

  • ❌ «Якщо виникла помилка, я просто роблю Thread.sleep(60000) прямо в циклі for-each» — Це гарантовано викличе перевищення max.poll.interval.ms, виліт консьюмера з групи і шторм ребалансувань.
  • ❌ «Kafka автоматично викидає биті повідомлення з топіка» — Kafka ніколи не змінює і не видаляє повідомлення з партиції. Логіка пропуску та ізоляції лежить на клієнті.
  • ❌ «При автокоміті (enable.auto.commit = true) дані ніколи не втрачаються» — Навпаки, автокоміт призводить до гарантованої втрати повідомлень при необроблених винятках.
  • ❌ «Паттерн Retry Topics можна застосовувати завжди без обмежень» — Він порушує суворий порядок повідомлень за ключем, що руйнівно для подієвих агрегатів.

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