📨 Раздел 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 можно применять всегда без ограничений» — Он нарушает строгий порядок сообщений по ключу, что разрушительно для событийных агрегатов.

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