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

如何在死信队列DLQ监听器中识别RabbitMQ错误处理器抛出的异常类型?

解决方案

核心思路是:异常信息不会自动传递到死信队列的消息头中,你可以在自定义错误处理器中手动把异常类型/定制化重试参数写入原生AMQP消息的头部,消息进入死信队列时会携带这些自定义头,在死信队列处理器中读取对应头值即可动态调整重试策略。


第一步:修改MyErrorHandler,写入自定义消息头

在抛出可重试异常前,把异常类型、定制化延迟值写入消息头,消息进入死信队列时会保留这些头信息:

public class MyErrorHandler implements RabbitListenerErrorHandler {

    // 自定义消息头key,避免和RabbitMQ原生头冲突
    public static final String X_EXCEPTION_TYPE = "x-exception-type";
    public static final String X_CUSTOM_RETRY_DELAY = "x-custom-retry-delay";

    @Override
    public Object handleError(Message amqpMessage,
            org.springframework.messaging.Message<?> message,
            ListenerExecutionFailedException exception) {

        Throwable cause = exception.getCause();
        // 检查异常是致命还是可重试
        if (cause instanceof FatalException || cause instanceof RetriesExceededException) {
            return new Status("FAIL!");
        }

        // 写入异常类型到消息头
        amqpMessage.getMessageProperties().getHeaders()
                .put(X_EXCEPTION_TYPE, cause.getClass().getSimpleName());
        // 限流异常直接写入更大的重试延迟,也可以统一放到DLQ中判断
        if (cause instanceof RateLimitException) {
            amqpMessage.getMessageProperties().getHeaders()
                    .put(X_CUSTOM_RETRY_DELAY, properties.getRateLimitRetryDelay());
        } else {
            amqpMessage.getMessageProperties().getHeaders()
                    .put(X_CUSTOM_RETRY_DELAY, properties.getRetryDelay());
        }

        // 抛出异常让消息NACK进入死信队列,此时消息已携带自定义头
        throw exception;
    }
}

第二步:修改MyDLQueueHandler,读取自定义头动态调整延迟

从死信消息的头中读取之前写入的异常标识/定制延迟值即可:

@Service
public class MyDLQueueHandler {

    private final MyAppProperties properties;
    private final MessageHeaderProcessor headerProcessor;
    private final RabbitProducerService rabbitProducerService;

    public MyDLQueueHandler(MyProperties properties, MessageHeaderProcessor headerProcessor, RabbitProducerService rabbitProducerService) {
        this.properties = properties;
        this.headerProcessor = headerProcessor;
        this.rabbitProducerService = rabbitProducerService;
    }

    /**
     * 由于仲裁队列不支持消息TTL,手动监听死信队列并带延迟重发消息,实现消息重试逻辑
     */
    @RabbitListener(queues = {"#{myDLQueue.getName()}"})
    public void handleError(@Payload Object message, @Headers MessageHeaders headers) {
        String routingKey = headerProcessor.getRoutingKey(headers);
        // 读取自定义延迟,无值则用默认配置兜底
        Integer retryDelay = headers.get(MyErrorHandler.X_CUSTOM_RETRY_DELAY, Integer.class);
        if (retryDelay == null) {
            retryDelay = properties.getRetryDelay();
        }
        // 也可以读取异常类型自行判断延迟,两种方式二选一:
        // String exceptionType = headers.get(MyErrorHandler.X_EXCEPTION_TYPE, String.class);
        // if (RateLimitException.class.getSimpleName().equals(exceptionType)) {
        //     retryDelay = properties.getRateLimitRetryDelay();
        // }
        Map<String, Object> newHeaders = Map.of(
                MessageHeaderProcessor.DELAY, retryDelay,
                MessageHeaderProcessor.RETRIES_HEADER, headerProcessor.getRetries(headers) + 1
        );
        rabbitProducerService.sendMessageDelayed(message, routingKey, newHeaders);
    }
}

你也可以根据业务需要扩展自定义头的内容,比如限流接口返回的官方建议等待时间、异常栈摘要等,只要是可序列化的类型都可以放入AMQP消息头传递。

内容的提问来源于stack exchange,提问作者Süleyman Gezsat

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.24 13:45:05