📨 Section 15 · Question #24

How to Handle Errors When Reading Messages

When consuming records from Apache Kafka, failures are inevitable: downstream databases time out, HTTP dependencies fail, JSON payloads contain corrupt bytes, or business rules...


🟢 Junior Level

The Core Problem of Consumer Error Handling

When consuming records from Apache Kafka, failures are inevitable: downstream databases time out, HTTP dependencies fail, JSON payloads contain corrupt bytes, or business rules are violated.

Unlike traditional message brokers (such as RabbitMQ or SQS) where an individual message can simply be negative-acknowledged (nack / reject) and returned to the queue head, a Kafka partition is an immutable, strictly ordered Append-Only Log. You cannot “pluck” a failed message out of the middle of a partition or skip past it without advancing the consumer group’s committed offset pointer.

Partition Log:
[Offset 100: OK] -> [Offset 101: ERROR (Poison Pill)] -> [Offset 102: OK] -> [Offset 103: OK]
                         │
                         ▼
        The Offset Commitment Dilemma:
        1. Commit Offset 101 -> SILENT DATA LOSS (Message 101 never successfully processed).
        2. Do not commit & crash -> INFINITE CRASH LOOP (On reboot, consumer re-reads 101 and crashes).

Core Scenarios and the Baseline Strategy

  1. Transient Errors: Temporary network blips, database deadlocks, or 503 gateway timeouts. Resolved via Retries with Exponential Backoff.
  2. Fatal Errors (“Poison Pills”): Corrupted binary bytes, incompatible schema deserialization failures, arithmetic errors. Retrying is futile; the message must be routed to a Dead Letter Queue (DLQ) for out-of-band analysis.

Safe Error Handling Pattern in Pure 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. Disable auto-commit to prevent committing failed offsets prematurely
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("Transient DB failure at offset {}, attempt {}", record.offset(), attempts);
                    Thread.sleep(1000); // Danger: must never exceed max.poll.interval.ms!
                } catch (Exception fatalEx) {
                    log.error("Fatal business error. Routing to DLQ: offset={}", record.offset(), fatalEx);
                    sendToDlqSynchronously("orders-dlq", record, fatalEx);
                    success = true; // Handled via DLQ; proceed forward
                    break;
                }
            }
            
            if (!success) {
                // Retries exhausted: isolate into DLQ
                sendToDlqSynchronously("orders-dlq", record, new RuntimeException("Retries exhausted"));
            }
        }
        
        // 2. Commit offsets ONLY after every batch record has processed or landed safely in DLQ
        consumer.commitSync();
    }
} catch (WakeupException e) {
    log.info("Consumer shutting down cleanly...");
} finally {
    consumer.close();
}

🟡 Middle Level

Error Classification and Remediation Matrix

| Error Category | Example Exceptions | Remediation Strategy | Impact on Message Order | | :— | :— | :— | :— | | Deserialization (Poison Pill) | SerializationException, malformed JSON/Avro | ErrorHandlingDeserializer + DLQ | Ordering preserved; bad record skipped to DLQ | | Transient Infrastructure Blip | SocketTimeoutException, LockAcquisitionException | In-line retry with backoff & pause | Strict FIFO order preserved; partition waits | | Business Logic Validation Failure | Negative account balance, missing required fields | Direct DLQ routing (Zero retries) | Malformed message isolated; pipeline proceeds | | Fatal Environment Crash | OutOfMemoryError, revoked database credentials | Immediate Container Shutdown | Application halts; alerts SRE immediately |


The max.poll.interval.ms Rebalance Trap

If a consumer implements naive blocking retries using Thread.sleep():

  1. The consumer thread blocks in a sleep loop and fails to invoke consumer.poll() within max.poll.interval.ms (default: 300,000 ms / 5 minutes).
  2. The Group Coordinator on the broker presumes the consumer thread has deadlocked or crashed.
  3. The coordinator triggers a Rebalance, stripping the partition from this consumer and assigning it to another healthy instance in the group.
  4. The new consumer polls from the last committed offset, encounters the exact same poison record, blocks in the same sleep loop, and times out as well.
  5. Result: A Cascading Rebalance Storm, halting all processing across the consumer group.

Spring Kafka Modern Architecture: DefaultErrorHandler

In Spring Boot applications, DefaultErrorHandler (introduced in Spring Kafka 2.8+ to replace SeekToCurrentErrorHandler) provides enterprise-grade error resilience:

@Configuration
public class KafkaConsumerConfig {

    @Bean
    public DefaultErrorHandler errorHandler(KafkaTemplate<Object, Object> template) {
        // 1. DLQ Publisher: Writes failed records to <originalTopic>.DLT
        DeadLetterPublishingRecoverer recoverer = new DeadLetterPublishingRecoverer(template,
            (record, exception) -> new TopicPartition(record.topic() + ".DLT", record.partition()));

        // 2. Exponential Backoff: 1s, 2s, 4s up to 3 attempts
        ExponentialBackOffWithMaxRetries backOff = new ExponentialBackOffWithMaxRetries(3);
        backOff.setInitialInterval(1000L);
        backOff.setMultiplier(2.0);
        backOff.setMaxInterval(10000L);

        DefaultErrorHandler errorHandler = new DefaultErrorHandler(recoverer, backOff);

        // 3. Route non-retryable exceptions to DLT immediately without delays
        errorHandler.addNotRetryableExceptions(
            IllegalArgumentException.class,
            JsonParseException.class,
            NoSuchElementException.class
        );

        return errorHandler;
    }
}

How DefaultErrorHandler Prevents Timeout Eviction

DefaultErrorHandler never sleeps the calling thread. Instead, when an exception occurs, it issues consumer.pause(partitions), calls consumer.seek() back to the failed offset, and executes rapid poll(Duration.ZERO) calls to maintain broker heartbeats. Once the backoff duration expires, it calls consumer.resume() and re-reads the failed record.


🔴 Senior Level

Non-Blocking Retries: The Delay Topics Pattern (@RetryableTopic)

In high-throughput microservices, pausing a partition for 5–10 minutes while waiting for a downstream service to recover causes unacceptable consumer lag spikes.

Modern architectures implement the Non-Blocking Delay Topics Pipeline (pioneered by Uber and standard in Spring Kafka):

[Main Topic: orders] ──► Consumer 1 (Downstream HTTP API Timeout)
                               │
                               ▼
        [Topic: orders-retry-10s] ──► Consumer 2 (Pulls with 10s delay)
                                           │
                                           ▼ (Fails again)
        [Topic: orders-retry-60s] ──► Consumer 3 (Pulls with 60s delay)
                                           │
                                           ▼ (All attempts exhausted)
        [Topic: orders-dlq]       ──► SRE Dashboard / Manual Investigation

Spring Kafka @RetryableTopic Implementation

@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("Processing order {} from topic {}", event.orderId(), topic);
        billingClient.charge(event);
    }

    @DltHandler
    public void processDlt(OrderEvent event, @Header(KafkaHeaders.ORIGINAL_OFFSET) long offset) {
        log.error("Order {} permanently failed after retries at offset {}. Sent to DLT.",
            event.orderId(), offset);
        alertingService.notifyOnCallSupport(event);
    }
}

[!CAUTION] The Critical Ordering Trade-Off: Forwarding failed records to separate retry topics breaks strict FIFO ordering for records sharing the same partition key. If an account produces a CREATE_ORDER event followed by a CANCEL_ORDER event, and the create event is diverted to a 60-second retry topic, the cancel event will process first on the main topic, corrupting state!


Deserialization Poison Pills: ErrorHandlingDeserializer

If a record contains corrupted bytes or an incompatible schema, standard deserializers (StringDeserializer, KafkaAvroDeserializer) throw a SerializationException inside consumer.poll().

  • Application code is never invoked.
  • The consumer catches the exception, rewinds, and immediately fails again on the next poll, deadlocking the partition indefinitely.

Architectural Solution: ErrorHandlingDeserializer

key.deserializer=org.apache.kafka.common.serialization.StringDeserializer
# Delegate deserialization through ErrorHandlingDeserializer
value.deserializer=org.springframework.kafka.support.serializer.ErrorHandlingDeserializer
spring.deserializer.value.delegate.class=org.springframework.kafka.support.serializer.JsonDeserializer

The ErrorHandlingDeserializer catches the parsing failure, supplies null as the payload value, and injects deserialization metadata headers:

  • DeserializationException.KEY_EXCEPTION_FQCN
  • DeserializationException.VALUE_EXCEPTION_STACKTRACE

Spring’s DefaultErrorHandler detects these headers and routes the raw corrupt bytes directly to the DLT without ever invoking the business listener!


4 Tricky Questions

1. How do you implement retries in vanilla Java KafkaConsumer without third-party frameworks and without risking max.poll.interval.ms eviction?

Answer: Use partition pausing and seeking:

  1. When a transient error occurs, record the failed TopicPartition and record offset.
  2. Call consumer.pause(Collections.singleton(topicPartition)).
  3. Reset the offset pointer: consumer.seek(topicPartition, failedOffset).
  4. In the primary application loop, continue invoking consumer.poll(Duration.ofMillis(100)) on schedule. Because the partition is paused, the broker returns zero records, but the poll execution continuously sends background heartbeats to the coordinator, preventing rebalance timeouts.
  5. Track elapsed time using an application backoff timer. When the backoff expires, call consumer.resume(Collections.singleton(topicPartition)).
  6. The next poll() re-fetches the failed record, preserving 100% strict partition FIFO ordering.

2. What happens if record #250 in a batch of 500 records fails while enable.auto.commit = true?

Answer: This results in Silent Data Loss:

  1. When enable.auto.commit = true, the background thread commits the highest fetched offset at auto.commit.interval.ms (default: 5s) during poll() calls.
  2. If poll() fetched offsets 1 to 500, and record 250 throws an exception that breaks out of the loop into the next poll() cycle, the consumer client will automatically commit offset 500 to __consumer_offsets.
  3. Records 250 through 500 remain uncompleted, yet their offsets are committed.
  4. If the application crashes or restarts, it resumes reading from offset 501. Records 250 through 500 are permanently lost.
  5. Rule: Robust production consumers must always configure enable.auto.commit = false.

3. Why must publishing to a Dead Letter Queue (DLQ) be synchronous or transactional prior to committing offsets?

Answer: If a consumer dispatches failed records to a DLQ asynchronously (dlqProducer.send(record)) without blocking:

  1. send() appends the message to the producer’s local in-memory RecordAccumulator buffer and immediately returns a Future.
  2. The consumer proceeds to invoke consumer.commitSync() on the main topic partition, believing the message has been secured.
  3. If the DLQ broker is unreachable or the message exceeds max.message.bytes, the producer’s background thread eventually fails and discards the record after exhausting retries.
  4. Because the original offset was already committed, the message is permanently gone from both the main topic and the DLQ.
  5. Invariant: The consumer must either wait for the DLQ ACK synchronously (dlqProducer.send().get()) or execute the DLQ produce and offset commit inside a unified Kafka Transaction (sendOffsetsToTransaction).

4. What are the major hazards of using @RetryableTopic for high-throughput topics with partition-key-sensitive logic?

Answer:

  1. Ordering Breakdown: Messages sharing the same business key are diverted to delay queues while newer messages on the main topic proceed uninterrupted, causing inverted state transitions.
  2. State Race Conditions: If event $E_1$ (Order Created) fails and lands in a 30-second delay topic, and event $E_2$ (Order Cancelled) processes successfully on the main topic 2 seconds later, $E_1$ will eventually execute 28 seconds later, resurrecting a cancelled order.
  3. Traffic Amplification: Each retry step requires producing and storing a new physical record across the broker cluster, multiplying network egress, broker disk I/O, and PageCache churn by the number of retry attempts.

🎯 Interview Cheat Sheet

Core Concepts

  • Append-Only Invariant: You cannot remove or skip an individual message in a Kafka partition; errors must be resolved via in-line retry, partition pause, or routing to a DLQ.
  • Poison Pills: Unparseable records that cause crash loops; resolved using ErrorHandlingDeserializer.
  • The Sleep Trap: Avoid Thread.sleep() in consumer threads; use consumer.pause() / seek() or Spring’s DefaultErrorHandler to preserve heartbeats.
  • Non-Blocking Retries: @RetryableTopic routes failed records to delay topics, trading FIFO ordering for zero head-of-line blocking.
  • Commit Safety: Disable enable.auto.commit and ensure DLQ publishes are confirmed synchronously before committing offsets.

Production Monitoring Metrics

| Metric | Meaning | Alert Threshold | | :— | :— | :— | | records-lag-max | Maximum consumer lag across assigned partitions | Sustained growth indicates stalled consumer | | app_kafka_consumer_dlq_total | Count of messages published to DLQ/DLT | $> 0$ triggers on-call investigation | | app_kafka_consumer_retry_rate | Proportion of consumed messages requiring retries | Spikes indicate downstream degradation | | last-heartbeat-seconds-ago | Elapsed time since last coordinator heartbeat | $> 0.5 \times \text{session.timeout.ms}$ |

Red Flags to Avoid

  • ❌ “If an error occurs, I just put Thread.sleep(60000) inside the listener loop.” (This triggers max.poll.interval.ms eviction and cascading rebalance storms).
  • ❌ “Kafka automatically evicts corrupt messages after 3 failed attempts.” (Kafka never modifies partition logs; error isolation is strictly the client’s responsibility).
  • ❌ “Auto-commit guarantees messages are never lost during errors.” (Auto-commit leads directly to silent data loss when exceptions occur mid-batch).
  • ❌ “Retry Topics should be used for every Kafka pipeline.” (Retry topics break message ordering; they cannot be used where entity state transitions depend on key ordering).