What does `enableDlq` do with the Kafka binder, and what happens to a poison message that exhausts its retries?
answer
- enableDlq=true → error.<destination>.<group>
- under kafka.bindings namespace, needs a group
- adds x-exception-* / x-original-* headers, commits original offset
- no DLQ = logged + offset committed = lost
- Rabbit analog = autoBindDlq / dead-letter exchange
basics
~20 sA poison message is one that keeps failing no matter how often you retry. With the Kafka binder's enableDlq=true, once retries are exhausted the message is sent to a separate dead-letter topic instead of being lost, so you can inspect and reprocess it later.
solid answer
~40 sSetting `spring.cloud.stream.kafka.bindings.<name>.consumer.enableDlq=true` tells the Kafka binder to route messages that still fail after all retries to a **dead-letter topic** rather than dropping them. The default DLQ topic name is `error.<destination>.<group>` (override with `dlqName`); the binding must have a `group`. The original payload and headers are preserved and the binder adds diagnostic headers like `x-original-topic`, `x-exception-message`, and `x-exception-stacktrace`. The main record's offset is then committed so the live consumer moves on — the poison message no longer blocks the partition. DLQ decouples *keep processing good traffic* from *don't silently lose bad messages*: you monitor the DLQ, fix the bug or bad data, and replay. Without `enableDlq`, an exhausted message on Kafka is logged and its offset committed — effectively discarded. RabbitMQ has an analogous `autoBindDlq` using a dead-letter exchange.
code
java · 32 lines// Enable DLQ + rename it (application.yml reference):
//
// spring.cloud.stream.bindings.process-in-0.destination: orders
// spring.cloud.stream.bindings.process-in-0.group: order-workers
// spring.cloud.stream.bindings.process-in-0.consumer.max-attempts: 3
// spring.cloud.stream.kafka.bindings.process-in-0.consumer:
// enable-dlq: true
// dlq-name: orders.dlq # default would be error.orders.order-workers
// dlq-partitions: 1
import java.util.function.Consumer;
import org.springframework.context.annotation.Bean;
import org.springframework.messaging.Message;
public class DlqExample {
// Main consumer: after 3 failed attempts, the record lands in orders.dlq
@Bean
public Consumer<Order> process(OrderService svc) {
return svc::handle;
}
// Optional: a monitor/replayer bound to the DLQ topic reads the diagnostic headers.
@Bean
public Consumer<Message<byte[]>> parkingLot() {
return msg -> {
Object origTopic = msg.getHeaders().get("x-original-topic");
Object reason = msg.getHeaders().get("x-exception-message");
System.out.printf("DLQ msg from %s failed: %s%n", origTopic, reason);
};
}
}go deeper
Know DLQ = a side topic for messages that keep failing, so they aren't lost.
State the default name error.<dest>.<group>, the group requirement, and the kafka.* namespace.
Describe preserved headers, offset commit/unblocking, dlqName/dlqPartitions, and the replay workflow.
Address best-effort (non-transactional) DLQ delivery, parking-lot governance, and DLQ depth alerting/SLOs.
**Poison message.** A message that *cannot* be processed successfully no matter how many times you retry — malformed payload, a schema your code can't parse, a business invariant it violates. Retrying it forever is pointless and, because retry is blocking, it stalls everything behind it. You need to get it *out of the way* without losing it. **Dead-letter queue (DLQ).** A separate destination where failed messages are parked for later inspection/replay. In Spring Cloud Stream this is a **binder-specific** feature. **Kafka binder.** Enable via `spring.cloud.stream.kafka.bindings.<bindingName>.consumer.enableDlq=true` (note the `kafka` segment — it's a *Kafka-specific* consumer property, not the generic binding namespace). Behavior: - After the binding-level retry (`maxAttempts`) is **exhausted**, the failing record is **published to a dead-letter topic**. - **Default topic name:** `error.<destination>.<group>`. So destination `orders`, group `order-workers` → `error.orders.order-workers`. Override with **`dlqName`**. - A **`group` is required** — the default name embeds it, and DLQ is a durable-group feature. - The original **key, payload, and headers are preserved**, and the binder adds diagnostic headers: `x-original-topic`, `x-original-partition`, `x-original-offset`, `x-exception-fqcn`, `x-exception-message`, `x-exception-stacktrace`. - The original record's **offset is committed**, so the live consumer advances — the partition is unblocked. - DLQ partitioning: by default the binder tries to send to the same partition number; `dlqPartitions` controls this (e.g. set to 1 to funnel all into a single partition). **RabbitMQ binder.** The analog is `spring.cloud.stream.rabbit.bindings.<name>.consumer.autoBindDlq=true`, which binds a **dead-letter exchange** and a `<destination>.<group>.dlq` queue; `republishToDlq` (default true) republishes with exception headers. **Without a DLQ.** On Kafka, an exhausted message is logged at ERROR and its offset committed → the message is **effectively lost** (no replay). This is the silent-data-loss trap teams hit in production. **Operational pattern.** DLQ isn't the end — it's a workflow: (1) alert on DLQ depth, (2) inspect the exception headers to diagnose, (3) fix code or data, (4) **replay** from the DLQ back to the main topic (a small re-publisher app or a `Function` reading `error.orders.order-workers` and re-emitting). Some teams add a *parking-lot* / manual-review step. **Gotchas.** - `enableDlq` lives under the `kafka.bindings` namespace, *not* the generic `bindings` one — a common misconfiguration that silently does nothing. - Producing to the DLQ can itself fail (broker down); that failure is logged, and the message may be lost — DLQ is best-effort, not transactional with the source read. - The DLQ topic must exist / be auto-createable; check `auto-create-topics`. - Don't treat DLQ as retry: infinite auto-replay of a truly poison message just loops. Gate replay behind a fix. **When to use.** Turn on `enableDlq` for any durable consumer where losing a message is unacceptable and where poison messages are plausible (external/unvalidated input). Pair it with non-retryable classification so known-bad messages reach the DLQ fast instead of burning retries.
- The DLQ config seems ignored — messages just vanish after retries. What's the most likely mistake?`enableDlq` was placed under the generic `spring.cloud.stream.bindings.<name>.consumer` namespace instead of the binder-specific `spring.cloud.stream.kafka.bindings.<name>.consumer`. Only the Kafka namespace enables the DLQ, so it silently did nothing and the message was dropped.
- How do you get a message out of the DLQ back into normal processing?Fix the underlying code/data, then run a replay: a small consumer bound to the DLQ topic that re-publishes to the original destination (reading `x-original-topic` from headers). Gate it behind the fix so you don't loop a still-poison message.
saying these in an interview costs you the question
- Believing DLQ routing works without a `group` set
- Putting enableDlq under the generic bindings namespace instead of kafka.bindings
- Assuming DLQ delivery is transactional/guaranteed with the source read (it's best-effort)
- Auto-replaying the DLQ in a loop without fixing the cause