How would you choose and design a Kafka error/retry/DLT strategy for a service with mixed ordering and throughput requirements?
answer
- classify transient vs fatal first
- ordering→blocking, throughput→non-blocking, soft→hybrid
- idempotent consumers are mandatory
- DLT: fail-safe publish + alert on depth + replay tool
- poison pill: ErrorHandlingDeserializer + bounded backoff
basics
~20 sMatch the mechanism to the constraint: use blocking seek retries (DefaultErrorHandler) where per-partition ordering matters and outages are short; use non-blocking @RetryableTopic where throughput and long back-offs matter and ordering can be relaxed. Always classify transient vs fatal, and route the unrecoverable to a DLT with replay tooling.
solid answer
~50 sI start from constraints, not tools. First, classify failures: **transient** (downstream timeout, lock contention, rate limit) vs **fatal/deterministic** (bad payload, deserialization, validation) — fatal ones must skip retries and go straight to the DLT. Second, the ordering-vs-throughput axis. Where strict per-key ordering is required (aggregate state machines), I use **blocking** retries (DefaultErrorHandler + ExponentialBackOff) so a stuck record holds its partition — accepting head-of-line blocking — and keep back-offs short to avoid rebalance risk and partition starvation. Where throughput dominates and reordering-on-failure is tolerable, I use **non-blocking** `@RetryableTopic` so the main stream never stalls, possibly with a small blocking retry first for micro-blips. Both terminate in a `DeadLetterPublishingRecoverer` DLT enriched with failure headers, with retention, alerting on DLT depth, and a replay job that reads the DLT, fixes/routes, and republishes to the source topic. I make the publish fail-safe (`failIfSendResultIsError`) so recovery can't silently lose data, and I instrument retries/DLT with metrics.
code
java · 22 lines// Ordering-critical stream: blocking, short exponential back-off, fatal->DLT now
@Bean
DefaultErrorHandler orderingCriticalHandler(KafkaTemplate<Object,Object> tmpl) {
var recoverer = new DeadLetterPublishingRecoverer(tmpl,
(rec, ex) -> new TopicPartition(rec.topic() + ".DLT", -1));
recoverer.setFailIfSendResultIsError(true); // never silently drop
var backOff = new ExponentialBackOffWithMaxRetries(4);
backOff.setInitialInterval(500);
backOff.setMaxInterval(5_000);
var h = new DefaultErrorHandler(recoverer, backOff);
h.addNotRetryableExceptions(ValidationException.class); // deterministic -> DLT
h.setRetryListeners((rec, ex, attempt) ->
meters.counter("kafka.retry", "topic", rec.topic()).increment());
return h;
}
// Throughput-critical, order-tolerant stream: non-blocking
@RetryableTopic(attempts = "5",
backoff = @Backoff(delay = 2_000, multiplier = 3.0, maxDelay = 60_000),
exclude = { ValidationException.class })
@KafkaListener(topics = "notifications")
void notify(Event e) { sender.send(e); } // must be idempotentgo deeper
Understand that different failures need different handling and unrecoverable ones go to a DLT.
Contrast blocking vs non-blocking and know fatal errors shouldn't be retried.
Reason from ordering/throughput constraints to a mechanism and include idempotency + DLT ops.
Own the end-to-end policy: failure taxonomy, per-stream blocking/non-blocking/hybrid choice, EOS/commit semantics, DLT topology, alerting, and replay tooling with an eye on cost and operability.
**Frame it as a decision, not a default.** The interviewer wants to see you reason from requirements to mechanism, and name the failure modes of each choice. **Step 1 — Classify failures (retryable vs fatal).** Build an explicit taxonomy: - *Transient/recoverable*: network timeouts, `503`s, DB deadlocks/`OptimisticLockException`, rate limits. Worth retrying (often with a real back-off). - *Fatal/deterministic*: `DeserializationException`, `MessageConversionException`, `ClassCastException`, validation/`IllegalArgumentException`, business-rule rejections. Retrying is pointless and, if blocking, actively harmful. Route straight to DLT. Map this with `addNotRetryableExceptions`/`addRetryableExceptions` (blocking) or `@RetryableTopic(include/exclude)` (non-blocking). Consider `setClassifications(map, false)` (whitelist) if most failures are deterministic. **Step 2 — Ordering vs throughput (the core axis).** - **Blocking (DefaultErrorHandler, seek-based):** preserves per-partition/per-key ordering because the failed record holds its offset and nothing behind it advances. Cost: head-of-line blocking — one poison-ish/slow record starves the whole partition; long back-offs risk exceeding `max.poll.interval.ms` (mitigated by partition pausing, but still throughput loss). Choose when correctness depends on order (event-sourced aggregates, CDC apply, ledgers). - **Non-blocking (@RetryableTopic):** main stream never stalls; supports long back-offs (minutes) cheaply. Cost: **loss of ordering** (retried records processed out of band), topic sprawl, more complex EOS. Choose for high-throughput, order-tolerant work (notifications, enrichment, idempotent upserts). - **Hybrid:** short in-place blocking retries (a couple of fast attempts) to absorb sub-second blips, then fall through to non-blocking retry topics for real outages — best of both when ordering is 'soft'. **Step 3 — Idempotency underpins everything.** Retries and reordering mean at-least-once and out-of-order delivery. Consumers must be **idempotent** (dedupe keys, upserts, conditional writes) — otherwise retries cause double side-effects. This is a precondition for non-blocking retries especially. **Step 4 — DLT design & operations.** Terminate every path in a `DeadLetterPublishingRecoverer` DLT: enriched headers (original topic/partition/offset, exception, stack trace), retention long enough to react, and a decoupled partition strategy (resolver → -1) to avoid partition-count coupling. Make recovery fail-safe: `setFailIfSendResultIsError(true)` + `setWaitForSendResultTimeout(...)` so a failed DLT publish fails recovery instead of dropping the record. Operationally: **alert on DLT lag/depth**, dashboard retry/DLT counters via `RetryListener`, and build a **replay tool** (read DLT → optionally transform/patch → republish to source or a reprocessing topic), guarding against replay storms. **Step 5 — Delivery/commit semantics.** Decide ack mode and whether you need EOS. With Kafka transactions, DLT publish + offset commit are atomic. Understand that blocking retries interact with `max.poll.interval.ms` and container pausing; size back-offs accordingly. **Step 6 — Poison-pill & backpressure protection.** A record that always fails but is marked retryable can loop; ensure a bounded BackOff and DLT so it exits. For deserialization poison pills, use an `ErrorHandlingDeserializer` so the container gets a `DeserializationException` (fatal) and DLPR can still carry the raw bytes to the DLT. **Decision summary.** - Strict order + short failures → **blocking** DefaultErrorHandler + ExponentialBackOff + DLT. - High throughput + order-tolerant + long back-off → **non-blocking** @RetryableTopic + DLT + @DltHandler. - Soft order + want to absorb blips cheaply → **hybrid** blocking-then-non-blocking. - Always: classify fatal→DLT-now, idempotent consumers, fail-safe DLT publish, DLT alerting + replay. **Common senior-trap answers to avoid:** 'just retry everything' (blocks partitions / retries poison pills), 'non-blocking is always better' (ignores ordering), 'DLT is enough' (without replay/alerting it's a graveyard), 'retries are free' (they cost throughput and can double side-effects without idempotency).
- Why is idempotency a prerequisite before you enable any retry strategy?Kafka retries give at-least-once delivery and, with non-blocking topics, out-of-order delivery. If the consumer isn't idempotent, a retried record re-applies side effects (double charge, duplicate email, double write). Idempotent handling (dedupe keys, upserts, conditional writes) makes retries safe.
- A team dumps everything into a DLT but never looks at it. What's wrong and what do you add?A DLT without alerting and replay is a data graveyard — failures accumulate unnoticed and are never recovered. Add monitoring/alerting on DLT depth and consumer lag, a @DltHandler or inspector, and a replay tool that transforms/fixes and republishes to the source or a reprocessing topic, with guards against replay storms.
- When would a hybrid blocking-then-non-blocking retry be the right call?When ordering is 'soft' — you'd like to keep order for the common case but can tolerate reordering for genuinely failing records. A couple of fast in-place blocking retries absorb sub-second blips without a topic hop; longer outages fall through to non-blocking retry topics so the partition doesn't stall.
saying these in an interview costs you the question
- 'Non-blocking retries are strictly better' — ignores the ordering guarantee you give up
- 'Just retry everything a lot' — blocks partitions, loops poison pills, doubles side-effects without idempotency
- 'A DLT means we're safe' — without replay/alerting it's an unmonitored graveyard
- Not mentioning idempotency at all when discussing retries
- Treating deserialization poison pills as retryable