skip to content

Design an end-to-end error-handling strategy for a partitioned Kafka consumer: retry, DLQ, ordering, and delivery guarantees. What are the trade-offs?

level: principalimportance: should knowfreq 30%

answer

  1. classify: transient→retry, poison→DLQ (retryableExceptions/defaultRetryable)
  2. blocking retry = head-of-line block + rebalance risk → DefaultErrorHandler / retry topics
  3. at-least-once ⇒ idempotent handlers
  4. DLQ best-effort, not transactional; monitor depth + gated replay
  5. partition order vs liveness trade-off; avoid skew

basics

~20 s

Classify errors: retry transient ones with bounded backoff, send permanently-failing (poison) ones to a DLQ. On Kafka, remember partition ordering means blocking retry stalls a whole partition, delivery is at-least-once (so be idempotent), and DLQ delivery is best-effort.

solid answer

~50 s

I separate **transient** from **permanent** failures. Transient (timeouts, brief outages) get **bounded retry with backoff**; permanent/poison (bad payloads, validation) are marked non-retryable so they skip straight to a **DLQ** (`enableDlq`, topic `error.<dest>.<group>`) with diagnostic headers. Because the built-in RetryTemplate is **blocking and per-partition**, long backoff creates head-of-line blocking and can breach `max.poll.interval.ms`, causing a rebalance — so for high throughput I switch to **non-blocking** retry via a `DefaultErrorHandler` (or dead-letter retry topics) so bad messages don't stall good ones. Kafka gives **at-least-once** delivery (offset committed after handling), so every handler must be **idempotent**. DLQ publishing is best-effort, not transactional with the source read, so I monitor DLQ depth, preserve original headers for diagnosis, and build a gated **replay** path. Partitioning preserves per-key order; I choose keys to avoid skew and accept that DLQ'd messages break strict ordering for that key.

code

java · 27 lines
java
import org.apache.kafka.clients.consumer.ConsumerRecord;
import org.springframework.cloud.stream.binder.kafka.ListenerContainerCustomizer;
import org.springframework.context.annotation.Bean;
import org.springframework.kafka.listener.AbstractMessageListenerContainer;
import org.springframework.kafka.listener.DeadLetterPublishingRecoverer;
import org.springframework.kafka.listener.DefaultErrorHandler;
import org.springframework.util.backoff.ExponentialBackOff;

public class ResilientConsumerConfig {

    // Non-blocking error handling: disable the binding RetryTemplate (max-attempts=1),
    // then install a DefaultErrorHandler with bounded exponential backoff that recovers
    // poison records to a DLT instead of blocking the partition indefinitely.
    @Bean
    public ListenerContainerCustomizer<AbstractMessageListenerContainer<?, ?>> customizer(
            DeadLetterPublishingRecoverer dlpr) {
        return (container, dest, group) -> {
            ExponentialBackOff backOff = new ExponentialBackOff(500L, 2.0);
            backOff.setMaxInterval(5_000L);
            backOff.setMaxElapsedTime(30_000L); // bound total time < max.poll.interval.ms
            DefaultErrorHandler handler = new DefaultErrorHandler(dlpr, backOff);
            // Never retry known-permanent errors -> straight to DLT:
            handler.addNotRetryableExceptions(IllegalArgumentException.class);
            container.setCommonErrorHandler(handler);
        };
    }
}

go deeper

for a junior

Not expected at this depth; know retry-then-DLQ exists.

for a middle

Can wire enableDlq + retry properties and note idempotency, but may miss the systemic trade-offs.

for a senior

Explains blocking retry's partition impact, at-least-once idempotency, and DLQ replay as a workflow.

for a principal

Frames the whole design as SLA-driven lever choices (blocking vs non-blocking, at-least- vs exactly-once, order vs liveness, best-effort DLQ vs transactions) and can justify each pick.

This is a systems-design answer weaving together every mechanism in the leaf. Structure it as **classify → retry → dead-letter → guarantees → ordering → operate**. **1. Classify failures.** Not all errors are equal: - **Transient**: downstream DB/HTTP timeout, broker hiccup, optimistic-lock clash. Retrying *helps*. - **Permanent / poison**: unparseable payload, schema mismatch, business-rule violation. Retrying *never* helps and, since retry blocks, actively harms. Encode this with `consumer.retryableExceptions` (map of type→retry) and `defaultRetryable=false` so only whitelisted transient types retry; everything else short-circuits to the DLQ. **2. Retry with backoff.** The binding-level `RetryTemplate` (`maxAttempts`, `backOffInitialInterval`, `backOffMultiplier`, `backOffMaxInterval`) handles transient faults. Critical property on Kafka: it is **blocking and in-band on the poll thread**, and consumption is **per partition**. So: - One slow-retrying message = **head-of-line blocking** for its entire partition (every key on it waits). - Cumulative backoff must stay under `max.poll.interval.ms` or Kafka evicts the consumer → **rebalance** + reprocessing. - Retry state is in-memory; a crash resets the count and Kafka redelivers. For high-throughput / long-backoff needs, **replace blocking retry** (`maxAttempts=1`) with a Spring Kafka **`DefaultErrorHandler`** installed via a `ListenerContainerCustomizer`, or adopt **non-blocking retry topics** (tiered `...-retry-0`, `-retry-1`, then DLT) so failing messages leave the main partition immediately and are retried out-of-band. **3. Dead-letter routing.** `spring.cloud.stream.kafka.bindings.<name>.consumer.enableDlq=true` routes exhausted/poison messages to `error.<destination>.<group>` (override `dlqName`, tune `dlqPartitions`). Headers `x-original-topic/partition/offset`, `x-exception-message/-stacktrace` capture provenance. The source offset is then committed so the pipeline keeps flowing. **Caveats:** DLQ publish is **best-effort** — if the broker is down when producing to the DLQ, that write can fail and the message can be lost; it's *not* atomic with the source read. For stronger guarantees you'd add Kafka **transactions**/exactly-once semantics, at a throughput cost. **4. Delivery guarantees & idempotency.** Spring Cloud Stream on Kafka is **at-least-once** by default: the offset is committed *after* successful handling, and redelivery happens on crash, rebalance, or retry. Therefore **every handler must be idempotent** — dedupe by message id / natural key, use upserts, or make side effects safe to repeat. Exactly-once is possible with Kafka transactions (`spring.cloud.stream.kafka.binder.transaction.*`) but adds latency and operational complexity; reserve it for money-movement-grade needs. **5. Ordering.** Partitioning (partition key → same partition → single consumer) buys **per-key ordering**. But error handling erodes it: sending a message to the DLQ while later same-key messages proceed means the DLQ'd one is now *out of order*; non-blocking retry topics likewise reorder. Decide whether strict per-key order or liveness matters more. Avoid **skew**: pick high-cardinality keys, or a custom `PartitionSelectorStrategy`; a hot partition throttles the whole key space and amplifies head-of-line blocking. **6. Operate.** DLQ is a *workflow*, not a graveyard: alert on **DLQ depth/rate**, dashboard the exception headers, provide a **gated replay** (a consumer on the DLQ that re-emits to the source once the bug/data is fixed — never an unconditional loop, or a truly poison message ping-pongs forever). Consider a **parking-lot** topic for manual review. Track redelivery/duplicate metrics to validate idempotency. **Trade-off summary.** - Blocking retry = simple, preserves order, but stalls partitions and risks rebalance. - Non-blocking retry topics = keeps main flow live and scalable, but reorders and adds topics/complexity. - DLQ = no silent loss + inspectability, but best-effort delivery and needs an operational replay process. - At-least-once = simple/robust but demands idempotency; exactly-once = stronger but slower/complex. - Fine-grained partition keys = ordering + parallelism, but risk of hot-partition skew. The principal-level answer names these levers *and* states which you'd pick for a given SLA (throughput vs strict ordering vs zero-loss) rather than reciting defaults.

  • Blocking retry preserves per-key order; non-blocking retry topics scale better but reorder. How do you choose?
    By SLA. If strict per-key ordering is a hard invariant (e.g. a state machine per account), keep blocking retry with short bounded backoff and accept reduced throughput. If liveness/throughput dominates and consumers are idempotent and order-tolerant, use non-blocking retry topics and handle reordering downstream (e.g. version/timestamp checks).
  • Your DLQ replay accidentally loops the same poison message forever. What went wrong and how do you prevent it?
    The replay re-emitted to the source unconditionally, so the still-unfixable message failed and returned to the DLQ endlessly. Gate replay behind a deployed fix, add an attempt counter/header to cap re-replays, and route exhausted messages to a manual parking-lot instead of auto-recycling.
  • Why isn't `enableDlq` enough to guarantee no message is ever lost?
    DLQ publishing is best-effort and not transactional with the source read/offset commit. If the broker is unavailable when writing to the DLQ, that write can fail and the message is lost. True no-loss needs Kafka transactions / exactly-once, trading throughput for the guarantee.

saying these in an interview costs you the question

  • Treating blocking retry as free — ignoring head-of-line blocking and rebalance risk
  • Assuming DLQ makes delivery transactional / guaranteed
  • Designing at-least-once pipelines without idempotent handlers
  • Claiming DLQ + retry topics preserve strict per-key ordering
  • Auto-replaying the DLQ without a gate on the fix

context