Overview
A message that a consumer cannot process is one of the most operationally dangerous states in a distributed system — not because the failure is catastrophic in itself, but because naive handling of it makes it catastrophic. The default behaviour for most messaging consumers without explicit failure handling is to retry. Retry again. Retry forever. The consumer stalls on the bad message; consumer lag grows; the downstream system that depended on this topic falls behind; eventually the lag exceeds the retention window and unconsumed messages are lost. One malformed payment message has blocked a queue that ordinarily clears 50,000 transactions per hour.
This is the poison message problem. The term comes from IBM MQ, where a message that causes a consumer to fail repeatedly is literally poisonous — it kills consumer throughput without any indication to the monitoring layer that a structural failure has occurred. The dead letter queue (DLQ) is the mechanism that removes the poison from the processing pipeline before it kills throughput, preserves the message for analysis and replay, and notifies the operations team that intervention is required.
The critical framing: the DLQ is a safety valve, not a solution. A message arriving in the DLQ is a symptom. The root cause is either a bug in the consumer, a schema mismatch between producer and consumer, or a permanent failure in a downstream dependency. The DLQ gives you time to diagnose and fix that root cause without losing the message or blocking the pipeline. It does not fix anything on its own.
Every message that lands in a DLQ must be actively claimed by a handler process — a DLQ consumer that classifies, alerts, and either replays or archives the message. A DLQ with no consumer is worse than no DLQ: you know messages are failing, but you are not acting on them, and the message retention clock is ticking. Wire the DLQ consumer before the first message arrives in it.
DLQ Topology Patterns
Two main DLQ topologies exist for Kafka deployments. Choosing between them is a trade-off between operational simplicity and observability granularity.
Per-topic DLQ: each source topic has a corresponding dead letter topic, typically named with a -dlt suffix — so payments.events.v1 has a dead letter topic payments.events.v1-dlt. This is the pattern Spring Kafka’s @RetryableTopic implements by convention. Advantages: trivial to determine which source topic a dead-lettered message came from; per-topic DLQ depth is a direct proxy for the health of the corresponding consumer. Disadvantage: topic count grows proportionally with source topic count, which matters in large Kafka estates.
Shared DLQ: all consumer groups route failed messages to a single dead.letter.queue topic, with metadata headers carrying the original topic, partition, offset, and consumer group. Advantages: centralised observability; one DLQ handler handles all failure types. Disadvantage: the DLQ handler must deserialise heterogeneous message types — which requires a schema registry and careful handling of schema version differences between message types.
For a regulated financial platform the recommendation is per-topic DLQ. The operational clarity of knowing that a growing payments.pacs008.v1-dlt depth means the pacs.008 consumer is in trouble — not some other consumer — outweighs the topic count overhead. In IBM MQ terminology, the DLQ corresponds to the queue manager’s DEADQ attribute, which is a single shared dead letter queue by design; MQ does not support per-queue DLQs natively, so the classification burden falls on the DLQ handler.
| Dimension | Kafka DLQ (-dlt) | MQ Dead Letter Queue (DEADQ) |
|---|---|---|
| Semantics | Dead letter is a regular Kafka topic; partitioned, replicated, log-retained | MQ local queue; single-depth FIFO; no partitioning |
| Config mechanism | @RetryableTopic annotation or DeadLetterPublishingRecoverer bean | DEADQ attribute on the queue manager via DEFINE QMGR or ALTER QMGR |
| Retention | Configurable log retention (hours or bytes) per DLT topic | Message expiry (EXPIRY attribute on the DLQ); default unlimited |
| Replay mechanism | DLQ consumer reads, fixes, republishes to source topic; or offset reset on source topic | AMQMDLQ handler republishes to original destination; or manual MQ Explorer move |
| Monitoring | Consumer lag on -dlt consumer group via Kafka consumer group offsets; Prometheus kafka_consumer_group_lag | MQ queue depth via ibmmq_queue_depth Prometheus exporter; MQSC DISPLAY QSTATUS |
Kafka DLQ Implementation
Spring Kafka 3.x’s @RetryableTopic annotation is the idiomatic way to configure Kafka DLQ behaviour in a Spring Boot application. It wires the retry topic chain, the backoff policy, and the dead letter topic creation automatically when the application starts.
package info.saib.payments.consumer;
import org.springframework.kafka.annotation.KafkaListener;
import org.springframework.kafka.annotation.RetryableTopic;
import org.springframework.kafka.retrytopic.DltStrategy;
import org.springframework.kafka.retrytopic.TopicSuffixingStrategy;
import org.springframework.retry.annotation.Backoff;
import org.apache.kafka.clients.consumer.ConsumerRecord;
@Component
public class PaymentEventConsumer {
// Three attempts total: 1 original + 2 retries
// Retry topics: payments.events.v1-retry-0, payments.events.v1-retry-1
// Dead letter: payments.events.v1-dlt
@RetryableTopic(
attempts = "3",
backoff = @Backoff(delay = 1000, multiplier = 2.0, random = true),
autoCreateTopics = "false", // pre-provision topics via Strimzi KafkaTopic CR
dltStrategy = DltStrategy.FAIL_ON_ERROR, // DLT publish failure is fatal, not silent
topicSuffixingStrategy = TopicSuffixingStrategy.SUFFIX_WITH_INDEX_VALUE,
kafkaTemplate = "kafkaTemplate"
)
@KafkaListener(topics = "payments.events.v1", groupId = "payments-processor")
public void consume(ConsumerRecord<String, PaymentEvent> record) {
paymentProcessingService.process(record.value(), record.headers());
}
// DLT handler: receives messages after all retry attempts are exhausted
@DltHandler
public void handleDlt(ConsumerRecord<String, PaymentEvent> record,
@Header(KafkaHeaders.EXCEPTION_CAUSE_FQCN) String exceptionClass,
@Header(KafkaHeaders.EXCEPTION_MESSAGE) String exceptionMsg) {
dlqHandlerService.classify(record, exceptionClass, exceptionMsg);
// classification: schema failure → alert + archive
// transient downstream → alert + schedule replay
// business rule violation → alert + manual review queue
}
}
spring:
kafka:
consumer:
group-id: payments-processor
auto-offset-reset: earliest
enable-auto-commit: false # manual offset commit; commit only on success or DLT publish
properties:
isolation.level: read_committed
max.poll.records: 50
max.poll.interval.ms: 30000 # tight bound — prevents broker treating consumer as dead during slow processing
listener:
ack-mode: MANUAL_IMMEDIATE
missing-topics-fatal: true # fail fast if DLT topic not pre-provisioned
# Strimzi KafkaTopic CR for DLT (apply separately via GitOps)
# metadata.name: payments.events.v1-dlt
# spec.partitions: same as source topic (3)
# spec.config:
# retention.ms: 604800000 # 7 days — must exceed your investigation SLA
# cleanup.policy: delete
The DLT consumer group naming convention: use payments-processor-dlt for the DLT consumer, distinct from payments-processor for the source topic. This keeps consumer lag metrics separated and allows alerting on DLT consumer lag independently of source topic lag.
IBM MQ Dead Letter Queue
IBM MQ’s DLQ is a first-class concept in the queue manager. Every queue manager has a single dead letter queue configured via the DEADQ attribute. Messages are moved to the DLQ by the queue manager itself when they cannot be delivered — for example, when the destination queue is full, the message has expired, or the message exceeds the maximum message length. Consumer-driven failure (a message the application cannot process) requires the application to explicitly move the message to the DLQ after exceeding its retry threshold.
# IBM MQ Dead Letter Queue Handler (AMQMDLQ) configuration
# Run as: runmqdlq SAIB.DEAD.LETTER.QUEUE SAIBQM1 < AMQMDLQ.conf
# --- Global parameters ---
INPUTQ(SAIB.DEAD.LETTER.QUEUE) # DLQ name, must match ALTER QMGR DEADQ()
INPUTQM(SAIBQM1) # queue manager name
WAIT(YES) # wait for new messages; don't exit when queue empty
RETRYINT(60) # retry interval in seconds for action(retry)
# --- Rules (evaluated in order; first match wins) ---
# Rule 1: ISO 20022 payment messages → regulatory hold queue, never auto-delete
REASON(MQRC_NOT_AUTHORIZED)
FORMAT(MQSTR)
APPLNAME(IPS*)
ACTION(FWD) # forward — do not discard
FWDQ(SAIB.PAYMENT.REGULATORY.HOLD)
FWDQM(SAIBQM1)
# Rule 2: Queue-full condition — retry up to 5 times before forwarding to ops queue
REASON(MQRC_Q_FULL)
RETRY(5)
ACTION(RETRY)
# Rule 3: Expired messages — log reason code and discard (expiry is by design)
REASON(MQRC_EXPIRY_ERROR)
ACTION(DISCARD)
# Rule 4: All other reasons → forward to ops monitoring queue
REASON(*)
ACTION(FWD)
FWDQ(SAIB.OPS.DEAD.MESSAGES)
FWDQM(SAIBQM1)
MQ DLQ reason codes to monitor explicitly: MQRC_Q_FULL (2053) indicates a backed-up destination; MQRC_NOT_AUTHORIZED (2035) indicates an ACL misconfiguration; MQRC_MSG_TOO_BIG_FOR_Q (2030) indicates a message size limit violation that requires a schema or application change. Each reason code has a different remediation path and should trigger a different alert severity.
Configure the DLQ itself with an explicit MAXMSGL larger than the largest legitimate message expected on any queue in the manager — a DLQ that rejects an oversized dead letter because its own MAXMSGL is too small silently discards the message. Set MAXMSGL(104857600) (100 MB) on the DLQ queue definition.
Poison Message Classification
Not all dead-lettered messages are equally dangerous or equally urgent. A DLQ handler that treats all failures identically — alert operations, wait for manual review — will generate alert fatigue and slow down response to the failures that actually matter. Classification at the DLQ handler separates the response.
Three classification categories cover the vast majority of production DLQ messages:
- Schema validation failure. The message does not conform to the expected Avro/Protobuf/JSON schema — wrong schema version, missing required field, unexpected field type. Root cause: producer schema version ahead of consumer schema version, or a producer deploying a breaking schema change without a migration strategy. Resolution: fix the consumer schema or repair the message and replay.
- Business rule violation. The message is structurally valid but contains a value that violates a business constraint — a payment amount that exceeds the customer’s limit, a creditor IBAN in a sanctioned country, a transaction timestamp in the future. Root cause: missing validation in the producer, or a change in business rules that the producer has not adopted yet. Resolution: fix the source system; replay only after confirming the same message would now be accepted.
- Downstream unavailability (transient vs permanent). The consumer attempted to call a downstream service — a core banking REST endpoint, an IBAN registry — that returned a transient error (503, timeout). If the downstream recovers before retry exhaustion, the message processes normally. If retry exhaustion happens while the downstream is still unavailable, the message lands in the DLQ. Root cause: downstream outage or consumer retry budget misconfigured too short for the expected outage duration. Resolution: once downstream recovers, replay the DLQ messages.
Event Replay Patterns
Replay is not a single operation — it is a family of patterns with different use cases, cost profiles, and risk levels. Choosing the right replay pattern for the situation avoids creating new problems while solving the original one.
DLQ consumer republish. The standard replay path: a DLQ consumer reads messages from the DLT topic, optionally transforms or repairs them, and publishes them back to the original source topic. The consumer must set Kafka headers on republished messages to track replay provenance. Spring Kafka’s DeadLetterPublishingRecoverer sets kafka_original_offset, kafka_original_partition, and kafka_original_topic automatically; add custom headers for X-Retry-Count, X-Original-Topic, and X-DLQ-Reason to enable idempotency checks in the downstream consumer.
Offset reset (full topic replay). When a consumer group has a systematic processing bug — the bug has been fixed and all messages processed in the affected window must be reprocessed — offset reset to the earliest offset or to a specific timestamp is faster than replaying from the DLT topic. This is a high-risk operation: it reprocesses all messages in the range, including ones that were processed successfully. Idempotency in the consumer is not optional for this approach; it is mandatory.
package info.saib.integration.dlq;
import org.apache.kafka.clients.consumer.ConsumerRecord;
import org.apache.kafka.common.header.Headers;
import org.springframework.kafka.annotation.KafkaListener;
import org.springframework.kafka.core.KafkaTemplate;
import org.springframework.kafka.support.KafkaHeaders;
@Service
public class DlqReplayService {
private final KafkaTemplate<String, byte[]> kafkaTemplate;
private final DlqClassifier classifier;
private final ReplayIdempotencyStore idempotencyStore;
// Consumes from DLT; only active when replay gate is open (circuit breaker pattern)
@KafkaListener(topics = "payments.events.v1-dlt", groupId = "payments-processor-dlt",
containerFactory = "dltListenerContainerFactory")
public void handleDlt(ConsumerRecord<String, byte[]> record) {
DlqClassification cls = classifier.classify(record);
switch (cls.getCategory()) {
case TRANSIENT_DOWNSTREAM:
if (cls.isDownstreamHealthy()) {
replay(record, cls);
} else {
// downstream still unhealthy: leave in DLT, alert ops
alertOps(record, cls);
}
break;
case SCHEMA_FAILURE:
case BUSINESS_RULE_VIOLATION:
// must not auto-replay; route to manual review
routeToManualReview(record, cls);
break;
case REGULATORY_HOLD:
// ISO 20022 payment: never auto-delete; forward to compliance queue
routeToComplianceHold(record, cls);
break;
}
}
private void replay(ConsumerRecord<String, byte[]> record, DlqClassification cls) {
String idempotencyKey = record.key() + ":" + record.offset();
if (idempotencyStore.alreadyReplayed(idempotencyKey)) {
log.warn("Skipping duplicate replay for key={} offset={}", record.key(), record.offset());
return;
}
ProducerRecord<String, byte[]> replay = new ProducerRecord<>(
cls.getOriginalTopic(), record.partition(), record.key(), record.value()
);
// Propagate original headers and add replay tracking headers
record.headers().forEach(h -> replay.headers().add(h));
replay.headers().add("X-Retry-Count", incrementRetryCount(record.headers()));
replay.headers().add("X-Original-Topic", cls.getOriginalTopic().getBytes());
replay.headers().add("X-DLQ-Reason", cls.getReasonCode().getBytes());
replay.headers().add("X-Replayed-At", Instant.now().toString().getBytes());
kafkaTemplate.send(replay);
idempotencyStore.markReplayed(idempotencyKey);
}
}
A message replayed from the DLT to the source topic may land in a consumer that already processed it successfully on an earlier attempt — if the consumer processed it but the DLT publish happened before the offset was committed. Without an idempotency check (keyed on the message’s business key, not the Kafka offset), the consumer applies the business operation twice. For payment debits, double-processing is a customer-visible incident. Implement idempotency at the consumer using a processed_events table keyed on the message’s business identifier before enabling any replay path.
DLQ Observability
A DLQ that is not monitored is not a safety valve — it is a silent message graveyard. Three alert conditions on the DLQ are mandatory; everything else is dashboard-level visibility.
- DLT consumer lag growing. If the DLT consumer group is falling behind — consumer lag on
payments.events.v1-dltis increasing — the DLQ handler itself is failing to process dead letters. Alert immediately; this means the safety valve is blocked. Usekafka_consumer_group_lag{group="payments-processor-dlt"} > 50with a 1-minute evaluation window. - DLT write rate spike. A sudden increase in the rate of messages being published to the DLT indicates a consumer has started failing rapidly. Alert on
rate(kafka_topic_messages_in_total{topic="payments.events.v1-dlt"}[5m]) > 5. More than 5 messages per second arriving in the DLT warrants an immediate investigation of the source topic consumer. - MQ DLQ depth growing. For IBM MQ, alert when
ibmmq_queue_depth{queue="SAIB.DEAD.LETTER.QUEUE"} > 100sustained for 5 minutes. A depth spike followed by a plateau indicates a batch of failures that the AMQMDLQ handler is not clearing.
The Grafana dashboard should show, per topic or per queue: DLT message rate (messages/min), DLT consumer lag, time-since-last-DLT-consumer-commit, and a table of the last 20 DLT messages with their reason codes. This is the first screen the on-call engineer opens when a DLQ alert fires.
ISO 20022 Message Failure Patterns
ISO 20022 payment messages — pacs.008 credit transfers, pain.001 customer initiations — have failure patterns that are qualitatively different from general event failures because the consequence of mishandling them is regulatory, not just operational.
pacs.008 schema validation failure. A pacs.008 that fails XSD validation before reaching SAMA IPS must be quarantined, not retried. The fix requires identifying whether the schema version is wrong (a producer version mismatch) or the data is malformed (a source system data quality issue). Both require human review before the message can be repaired and resubmitted.
Mandatory field absent. SAMA IPS rejects pacs.008 messages missing SttlmMtd, ClrSys, or a Purp/Cd from the approved purpose code list. These are not schema validation failures — the message is XSD-valid but fails SAMA’s business rule layer. They arrive as pacs.002 reject messages with specific reason codes and must be routed back through the integration layer to the originating channel as a payment failure notification.
Routing DLQ with regulatory hold requirement. A payment message that fails after IPS acceptance — for example, an OFAC or UN sanctions hit that was not caught before submission — is not simply a failed message. It may be subject to a regulatory hold that requires preservation for compliance purposes. The DLQ handler for payment messages must never auto-delete or auto-expire messages that carry a sanctions or AML hold flag.
A pacs.008 that fails SAMA IPS processing and lands in the DLQ with a reason code indicating a suspected sanctions or AML match must be preserved in accordance with SAMA’s Technology Risk Management framework and the Saudi Anti-Money Laundering Law. Configure DLT retention for payment topics at a minimum of 90 days; for any message flagged with a compliance hold, archive to immutable cold storage before the DLT retention window expires. Never configure a Kafka cleanup.policy=compact or cleanup.policy=delete on a payment DLT without explicit compliance team sign-off.
Pitfalls
If a consumer bug sends messages to the DLT and you replay from the DLT before deploying the fix, the replayed messages hit the same consumer, fail again, and go back to the DLT. You now have a duplicate of every failed message in the DLT, and the next replay doubles the problem again. The replay gate — the mechanism that enables DLT consumer replay — must be disabled by default and enabled only after the consuming service is confirmed healthy at the new version in the target environment.
If your DLT topic has a 24-hour retention policy and your on-call SLA allows up to 48 hours before a DLQ alert must be investigated, every dead-lettered message will be expired before it can be examined. Align DLT retention with your P2 incident investigation SLA, not with the source topic retention. A minimum of 7 days is reasonable for most integration workloads; payment messages warrant 90 days.
- No DLQ on the DLQ handler. The DLQ handler is itself a Kafka consumer. If it fails processing a DLT message, and it has no error handling of its own, it will block on that message indefinitely. The DLQ handler needs its own retry policy — typically a simple in-memory retry with backoff, then an alert and skip, since a DLQ handler must never itself dead-letter to another topic (infinite poison).
- Replay order inversion breaking idempotency. When replaying multiple messages from the DLT, Kafka does not guarantee that messages are replayed in the same order they were originally produced. If your consumer is idempotent by message key but order-dependent within a key — for example, a payment status update that must follow the original payment event — replay order inversion can cause incorrect state transitions. Use per-partition ordering in the DLT and replay partition-by-partition.
- Treating the DLQ as a permanent archive. Engineers who find the DLQ a convenient place to leave messages “for later” create technical debt that compounds. Every message in the DLQ represents an unresolved failure. Measure DLQ age (time between message dead-lettering and message resolution) as an operational metric; set a P3 SLA for DLQ message resolution (e.g. within 5 business days), not just a P1 SLA for the initial alert.
Production Checklist
- Introduce DLQ by configuring
@RetryableTopicwithautoCreateTopics=falseand pre-provisioning the DLT topics via StrimziKafkaTopicCRs in Git before deployment. - Set DLT retention at least equal to your P2 incident investigation SLA; use a minimum of 7 days for non-payment topics, 90 days for payment topics.
- Wire the DLQ handler service before enabling retry exhaustion on the source consumer — the DLT must have a consumer before the first message arrives in it.
- Implement idempotency in all consumers using a business-key-based
processed_eventstable before enabling any replay path from the DLT. - Enable the replay gate only after the consuming service is confirmed healthy at the fixed version; disable it again after replay completes.
- Add
X-Retry-Count,X-Original-Topic,X-DLQ-Reason, andX-Replayed-Atheaders to all replayed messages to enable downstream idempotency checks and audit tracing. - Alert on DLT consumer lag growing, DLT write rate spike, and MQ DLQ depth growing; link each alert to a runbook with the classification and replay procedure.
- Configure a Grafana dashboard showing per-topic DLT depth, DLT write rate, DLT consumer lag, and a table of recent DLT messages with reason codes.
- For IBM MQ, configure AMQMDLQ with payment-specific rules that route to a regulatory hold queue and never auto-discard; verify the rule set covers all expected
MQRC_reason codes. - Never configure
cleanup.policy=compacton a payment DLT topic; immutably archive any DLT message with a compliance or regulatory hold flag before DLT retention expiry. - Measure DLQ message age (time-to-resolution) as an operational SLI; set a P3 SLA for resolution, not just a P1 SLA for the alert.
- Test the DLQ handler with failure injection: schema validation failure, transient downstream 503, business rule violation, and regulatory hold scenarios in a lower environment before production rollout.