You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何在死信队列(DLQ)中自定义异常消息或添加自定义头部参数

Customizing DLQ Messages with Custom Exception Headers/Messages in Spring Cloud Stream

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.12 05:10:13