如何在死信队列(DLQ)中自定义异常消息或添加自定义头部参数
Alright, let's figure out how to add custom exception messages or headers to your Dead Letter Queue (DLQ) messages in Spring Cloud Stream. Right now, your DLQ only gets the raw stack trace in headers, so here are a few practical approaches to fix that:
Option 1: Extend DeadLetterPublishingRecoverer for Automatic DLQ Handling
If you want to keep the built-in retry mechanism and just enhance the DLQ messages with custom headers, extending the default DeadLetterPublishingRecoverer is a great fit. This lets you modify the message before it's sent to the DLQ.
First, create a configuration class to define your custom recoverer:
@Configuration public class CustomDlqConfig { @Bean public Consumer<Message<?>> dlqErrorHandler(KafkaTemplate<?, ?> kafkaTemplate, DestinationResolver destinationResolver) { // Extend the default recoverer to add custom headers/messages DeadLetterPublishingRecoverer customRecoverer = new DeadLetterPublishingRecoverer(kafkaTemplate, destinationResolver) { @Override protected Message<?> createDeadLetterMessage(Message<?> originalMessage, Throwable exception) { // Build a custom error message using the original payload and exception details String customErrorMsg = String.format("Failed to process message with ID: %s - Error: %s", ((Content) originalMessage.getPayload()).getData().getId(), exception.getMessage()); // Create a new message with custom headers + original content return MessageBuilder.fromMessage(originalMessage) .setHeader("X-Custom-Error-Message", customErrorMsg) .setHeader("X-Exception-Type", exception.getClass().getSimpleName()) // Optional: Keep the original stack trace header if needed .build(); } }; // Return a consumer that handles error channel messages return errorMessage -> { Throwable exception = (Throwable) errorMessage.getHeaders().get(MessageHeaders.ERROR_CHANNEL_EXCEPTION); customRecoverer.accept(errorMessage, exception); }; } }
Make sure your existing config keeps enableDlq: true—this approach works with the automatic retry/DLQ flow:
messageInputChannel: consumer: enableDlq: true dlqName: error.message backOffInitialInterval: 15000 backOffMaxInterval: 30000 maxAttempts: 3 headerMode: raw
Option 2: Manually Handle Exceptions and Send to DLQ
If you want full control over when and how messages go to the DLQ, you can catch exceptions directly in your consumer method and construct a custom message to send to the DLQ.
First, update your sink interface to include the DLQ output channel:
public interface CustomStreamBindings { String INPUT_CHANNEL = "messageInputChannel"; String DLQ_CHANNEL = "error.message"; @Input(INPUT_CHANNEL) SubscribableChannel input(); @Output(DLQ_CHANNEL) MessageChannel dlqOutput(); }
Then modify your consumer to handle exceptions and send custom messages:
@StreamListener(CustomStreamBindings.INPUT_CHANNEL) public void consumeMessage(Message<Content> message, @Output(CustomStreamBindings.DLQ_CHANNEL) MessageChannel dlqChannel) { try { myService.update(message.getPayload().getData().getId()); } catch (Exception e) { // Build your custom error message and headers String customErrorMsg = String.format("Processing failed for ID: %s - %s", message.getPayload().getData().getId(), e.getMessage()); Message<?> dlqMessage = MessageBuilder.fromMessage(message) .setHeader("X-Custom-Error", customErrorMsg) .setHeader("X-Exception-Details", e.getClass().getName()) .build(); // Send the custom message to DLQ dlqChannel.send(dlqMessage); } }
Don't forget to disable the automatic DLQ in your config to avoid duplicate messages:
messageInputChannel: consumer: enableDlq: false # Disable auto-DLQ since we're handling it manually backOffInitialInterval: 15000 backOffMaxInterval: 30000 maxAttempts: 3 headerMode: raw
Option 3: Listen to the Error Channel with @ServiceActivator
If you don't want to modify your consumer logic at all, you can listen to the input channel's error stream, modify the message, and forward it to the DLQ.
Add this service activator to your code:
@ServiceActivator(inputChannel = "messageInputChannel.errors") public void handleChannelErrors(Message<?> errorMessage, @Output("error.message") MessageChannel dlqChannel) { // Extract the original message and exception from the error message Message<?> originalMessage = (Message<?>) errorMessage.getPayload(); Throwable exception = (Throwable) errorMessage.getHeaders().get(MessageHeaders.ERROR_CHANNEL_EXCEPTION); // Create custom error details String customErrorMsg = String.format("Error processing message: %s - %s", originalMessage.getPayload(), exception.getMessage()); // Build and send the custom DLQ message Message<?> dlqMsg = MessageBuilder.fromMessage(originalMessage) .setHeader("X-Custom-Error-Message", customErrorMsg) .build(); dlqChannel.send(dlqMsg); }
Keep enableDlq: false in your config here too—we're taking over error handling from the default DLQ mechanism.
Which Option Should You Choose?
- Option 1: Best if you want to keep the automatic retry/DLQ flow and just add custom headers.
- Option 2: Ideal if you need fine-grained control over exception handling and DLQ logic.
- Option 3: Perfect if you don't want to touch your existing consumer code.
内容的提问来源于stack exchange,提问作者oreoman

