Ch. 26 · Apache Kafka

Kafka Poison Messages, Retries and Dead Letter Topics

Why one bad Kafka record can block a whole partition, and how to handle it with error classification, bounded retries and a dead letter topic.

~9 min readintermediateupdated Oct 6, 2026

“Your consumer keeps crashing on one malformed message. What do you do?” Interviewers ask this because it separates people who have run Kafka in production from people who have only read about it. The naive answers, “skip it” or “retry until it works”, each cause a different outage. The follow-ups probe the design: “how do you tell a bad message from a database outage?”, “what goes in a dead letter topic?”, “what happens to ordering when you retry later?” and “how do you replay the dead letters once the bug is fixed?”

This guide explains why poison messages stall a partition, how to classify failures, and how to build bounded retries and a dead letter topic (DLT) with the plain Java client, with notes on Spring Kafka and Kafka Connect. It targets Kafka 3.8/4.0. Code is illustrative: it follows the documented APIs but was not run against a cluster.

Before you start

You should understand offset commits and at-least-once processing; Kafka offset commits and delivery semantics covers them. You should also know what a serializer and deserializer do, and be comfortable with Java exceptions.

The short answer

A poison message is a record that fails every time it is processed, such as invalid JSON, an unknown schema or data that breaks a business rule. Because an at-least-once consumer only commits after success, it restarts at the same offset and fails again, so the whole partition is blocked and lag grows. The fix is to classify errors: catch deserialization failures and step past them, retry transient errors (timeouts, an unavailable database) with backoff, and after a bounded number of attempts send permanent failures, with their original bytes and error details, to a dead letter topic, then commit and move on. The DLT is then monitored, fixed and replayed deliberately.

How it works

A partition is processed in order by one consumer in the group. If record 1,042 throws and the code does not commit past it, three things can happen, all bad:

  • The exception escapes the poll loop, the consumer dies, the orchestrator restarts it, and it reads record 1,042 again: a crash loop.
  • The code retries forever in place: the partition makes no progress and, if retries take longer than max.poll.interval.ms, the member is kicked out and the next owner repeats it: a rebalance loop.
  • The code catches and ignores the exception, commits, and moves on: silent data loss.

Other partitions keep flowing in all three cases, so the symptom is usually lag growing on one partition while the rest are at zero.

Errors fall into three classes, and the response follows from the class:

Class Examples Response
Deserialization malformed JSON, unknown schema ID, wrong format will never succeed: send raw bytes to the DLT, skip
Transient timeout, connection refused, HTTP 503, lock timeout retry with backoff; it will likely succeed later
Permanent validation failure, missing reference data, a bug bounded retries at most, then the DLT

In the Java client, a deserializer that throws inside poll() surfaces as a RecordDeserializationException, which reports the partition and offset of the bad record so you can seek past it; Kafka 3.8 added accessors for the raw key and value bytes (KIP-1036). Many teams avoid the question entirely by consuming byte[] and deserializing inside their own loop, where every failure is an ordinary exception with the bytes in hand.

Step-by-step walkthrough

Step 1: Deserialize inside the loop

Consuming raw bytes means one bad record can never break poll() itself, and the original bytes are available for the DLT.

props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, ByteArrayDeserializer.class);
props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, ByteArrayDeserializer.class);
props.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, false);

for (ConsumerRecord<byte[], byte[]> raw : consumer.poll(Duration.ofMillis(500))) {
    OrderEvent event;
    try {
        event = mapper.readValue(raw.value(), OrderEvent.class);
    } catch (IOException e) {
        deadLetter(raw, e, 0);            // malformed: retrying cannot help
        continue;
    }
    handleWithRetry(raw, event);
}
consumer.commitSync();
java

Step 2: Retry transient errors with a small, bounded backoff

In-place retries preserve ordering, but they hold up the partition and count against max.poll.interval.ms, so keep the total short. Anything longer belongs in a retry topic or a pause.

void handleWithRetry(ConsumerRecord<byte[], byte[]> raw, OrderEvent event) {
    int attempt = 0;
    while (true) {
        try {
            orderService.apply(event);
            return;
        } catch (TransientException e) {                 // timeouts, 503s, connection errors
            if (++attempt >= 4) { deadLetter(raw, e, attempt); return; }
            sleep(Duration.ofMillis(200L * (1L << attempt)));   // 400, 800, 1600 ms
        } catch (RuntimeException e) {                   // validation errors and bugs
            deadLetter(raw, e, attempt + 1);
            return;
        }
    }
}
java

The classification is the important part. Retrying a validation error four times wastes time; dead-lettering a database timeout on the first attempt floods the DLT during every blip.

Step 3: Publish a useful dead letter

A dead letter must contain everything needed to understand and replay it: the original key and value bytes untouched, plus where it came from and why it failed. Headers keep that metadata out of the payload.

void deadLetter(ConsumerRecord<byte[], byte[]> raw, Exception e, int attempts) {
    ProducerRecord<byte[], byte[]> dlt = new ProducerRecord<>("orders.DLT", null, raw.key(), raw.value());
    Headers h = dlt.headers();
    raw.headers().forEach(h::add);                                  // keep tracing headers
    h.add("dlt.original.topic", raw.topic().getBytes(UTF_8));
    h.add("dlt.original.partition", Integer.toString(raw.partition()).getBytes(UTF_8));
    h.add("dlt.original.offset", Long.toString(raw.offset()).getBytes(UTF_8));
    h.add("dlt.exception", e.getClass().getName().getBytes(UTF_8));
    h.add("dlt.message", String.valueOf(e.getMessage()).getBytes(UTF_8));
    h.add("dlt.attempts", Integer.toString(attempts).getBytes(UTF_8));
    try {
        dltProducer.send(dlt).get();   // block: never commit past a record whose dead letter was not written
    } catch (InterruptedException | ExecutionException ex) {
        throw new IllegalStateException("dead letter write failed", ex);   // stops the loop before the commit
    }
}
java

Blocking here is deliberate. If the DLT write fails, the exception stops the loop before the commit, so the record is retried rather than lost. Give the DLT long retention, such as 14 to 30 days, so there is time to investigate.

Step 4: Choose blocking or non-blocking retries

For longer outages, a common design adds retry topics with delays, such as orders.retry.1m and orders.retry.10m, each read by a consumer that waits until the record’s due time. The main topic keeps flowing, but a failed event for order-1 is now processed after later events for order-1.

Approach Ordering per key Partition progress
Blocking retry in place preserved stalls during retries
Retry topics, then DLT broken for retried keys continues
Park the key: once a key fails, route its later records to the DLT too preserved for that key continues for other keys

Pick blocking retries for state changes that must apply in order, such as account balances, and retry topics for independent events, such as emails.

Step 5: Replay deliberately

After a fix ships, a small replay tool reads the DLT, filters by exception or time range, and republishes the original key and value to the source topic. Keeping the key keeps per-key routing. Replayed records must be safe to process twice, because some may have partially succeeded the first time.

Worked scenario

A pricing team deployed a producer change that sent amount as the string "12.50" instead of a number. The order consumer’s JSON mapping failed on the first such record in partition 3. The exception escaped the loop, the pod crashed, Kubernetes restarted it, and it failed on the same offset. Within 20 minutes partition 3 had 40,000 records of lag while the other 11 partitions were at zero, and the crash loop also triggered a rebalance every restart, slowing the healthy partitions.

// Broken: any exception escapes and kills the consumer
for (ConsumerRecord<String, OrderEvent> r : consumer.poll(Duration.ofMillis(500))) {
    orderService.apply(r.value());
}
consumer.commitSync();
java

With the steps above, the bad records would have been written to orders.DLT with the exception MismatchedInputException and their original bytes, the partition would have continued, and an alert on the DLT rate would have fired within a minute. The immediate fix was to deploy the classifying consumer, then add schema enforcement through a schema registry so the producer could not publish the incompatible change at all. Once the consumer accepted both formats, the team replayed the dead letters to the main topic.

Common mistake

  • “Catch everything, log it and commit.” That is silent data loss with a log line.
  • “Retry until it succeeds.” A permanent error never succeeds, so the partition stalls forever.
  • Sending the parsed object, or a re-serialized one, to the DLT. Keep the original bytes; the parser is what failed.
  • A DLT nobody watches. Without alerts and an owner, it is the same as discarding records.
  • Assuming retry topics preserve order. They do not for the keys that were retried.

Verify the behavior

Produce a valid and an invalid record, then check that the consumer keeps going and the dead letter carries its metadata:

kafka-console-producer.sh --bootstrap-server localhost:9092 --topic orders \
  --property parse.key=true --property key.separator=:
# order-1:{"orderId":"order-1","amount":12.5}
# order-2:{"orderId":"order-2","amount":
# order-3:{"orderId":"order-3","amount":7.0}

kafka-console-consumer.sh --bootstrap-server localhost:9092 --topic orders.DLT --from-beginning \
  --property print.key=true --property print.headers=true
# expect order-2 with dlt.original.offset and dlt.exception headers

kafka-consumer-groups.sh --bootstrap-server localhost:9092 --describe --group orders-service
# LAG should return to 0: order-3 was processed after order-2 failed
Terminal

In a unit test, feed the handler a malformed payload and assert that the DLT producer received one record and that processing continued with the next.

Follow-up questions

How does Spring Kafka handle this? DefaultErrorHandler with a backoff retries in place, then a DeadLetterPublishingRecoverer sends the record to <topic>.DLT by default; ErrorHandlingDeserializer wraps deserialization failures, and @RetryableTopic sets up non-blocking retry topics.

What does Kafka Connect offer? Sink connectors support errors.tolerance=all with errors.deadletterqueue.topic.name, and errors.deadletterqueue.context.headers.enable=true adds the failure context as headers.

How many partitions should a DLT have? Few are usually enough, because volume should be low. Keep the original key so replays route correctly.

How do you prevent poison messages at the source? Enforce schemas with a schema registry and compatibility rules, and validate in the producer before sending.

Interview exercise

A consumer processes account transactions keyed by account ID. Some records fail because a downstream fraud API times out for up to 10 minutes during incidents; others fail because the account does not exist yet. Design the error handling, and explain what happens to ordering for an account whose record failed.

Answer and reasoning

The fraud API timeout is transient, and the missing account is likely permanent or a sequencing problem, so they get different paths. Because these are balance-changing transactions, per-account order matters, so retry topics that let later transactions for the same account overtake a failed one are dangerous. For the timeout, a better choice is blocking retry with backoff using pause() and seek() back to the failed offset, so the consumer keeps polling and stays in the group while the partition waits. For the missing account, retry a few times briefly, then dead-letter and park the key: route later records for that account to the DLT as well, so they are not applied out of order, and alert. After the fix, replay that account’s records in their original order, relying on idempotent handling by transaction ID.

Continue learning

Practise with the Apache Kafka interview questions and the Kafka MCQs. Related notes: Kafka offset commits and delivery semantics, retry budgets in microservices and consumer deduplication. Primary sources: the KafkaConsumer Javadoc for Kafka 4.0, the Spring for Apache Kafka guide to non-blocking retries and DLTs and the Apache Kafka documentation.

More in Apache Kafka

esc