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

基于DLQ的可重试AMQP消息处理方案优化问询

免责声明!

以下代码可正常运行,仅咨询实现的正确性及其他可行解决方案。


定义
  1. some-domain-name.q.some-service.incoming-messages - 该队列存储所有待处理的入站消息。
    incoming-messages queue

  2. some-domain-name.q.some-service.incoming-messages.postponed - 该队列存储所有无法处理的延迟入站消息(如遇业务、连接或其他异常)。队列设置了TTL,到期后消息将被移回队列(1)。
    incoming-message-postponed queue

  3. some-domain-name.q.dead-letter - 主死信队列(DLQ),存储所有无法处理的消息。
    dlq

流程

消息进入队列(1) some-domain-name.q.some-service.incoming-messages,若处理时发生异常,将被移至队列(2) some-domain-name.q.some-service.incoming-messages.postponed,同时自定义头x-retry-counter递增。等待TTL(如30秒)后,消息将再次进入队列(1);若再次异常,计数器递增并移至队列(2)。设置了最大重试次数,若次数耗尽,消息将被移至队列(3) some-domain-name.q.dead-letter。

存在一个特殊场景:若抛出SkipIncrementRetryCounterException,流程正常执行,但x-retry-counter头不递增。

实现

AmqpListener

该注解提供了便捷的队列处理、错误处理方式,并通过@SendTo将重试次数耗尽的消息发送至主DLQ(3)。

@RabbitListener(errorHandler = "amqpRabbitListenerErrorHandler")
@SendTo("!{@dlqNameProvider.apply(request)}")
public @interface AmqpListener {
    @AliasFor(annotation = RabbitListener.class, attribute = "queues")
    String[] value() default {};
}

AmqpRabbitListenerErrorHandler

该处理器决定消息应重试(抛出异常)还是移至主DLQ(3)(返回消息)。

public class AmqpRabbitListenerErrorHandler implements RabbitListenerErrorHandler {
    @Override
    public Object handleError(Message amqpMessage, org.springframework.messaging.Message<?> message, ListenerExecutionFailedException exception) {
        MessageProperties messageProperties = amqpMessage.getMessageProperties();

        int maxRetries = 3;
        int retryCount = getCounter(messageProperties);
        boolean shouldRetry = retryCount < maxRetries;
        if (shouldRetry) {
            throw exception;
        }

        return message;
    }

    private static Integer getCounter(MessageProperties messageProperties) {
        Integer retryCount = messageProperties.getHeader("x-retry-counter");
        if (retryCount == null) {
            retryCount = 0;
        }
        return retryCount + 1;
    }
}

AmqpRepublishMessageRecoverer

MessageRecoverer负责根据需要递增重试计数器,并将消息发送至延迟队列(2)。

public class AmqpRepublishMessageRecoverer extends RepublishMessageRecoverer implements BeanFactoryAware {
    private static final SpelExpressionParser spelExpressionParser = new SpelExpressionParser();

    public AmqpRepublishMessageRecoverer(AmqpTemplate errorTemplate) {
        super(
                errorTemplate,
                spelExpressionParser.parseExpression("@errorExchangeProvider.apply(#this)"),
                spelExpressionParser.parseExpression("@errorRoutingKeyExpression.apply(#this)")
        );
    }

    @Override
    public void setBeanFactory(BeanFactory beanFactory) throws BeansException {
        ((StandardEvaluationContext) evaluationContext).setBeanResolver(new BeanFactoryResolver(beanFactory));
    }

    @Override
    protected Map<? extends String, ?> additionalHeaders(Message message, Throwable cause) {
        if (!(cause.getCause() instanceof SkipIncrementRetryCounterException)) {
            Integer retryCount = message.getMessageProperties().getHeader("x-retry-counter");
            retryCount = retryCount == null ? 1 : retryCount + 1;
            return Map.of("x-retry-counter", retryCount);
        }
        return emptyMap();
    }
}

dlqNameProvider、errorExchangeProvider和errorRoutingKeyExpression

这些辅助函数仅用于提供RabbitMQ目标对象的名称,通过amqp_consumerQueue头计算取值。

RabbitRetryTemplateCustomizer

当抛出SkipIncrementRetryCounterException时,使用以下重试和退避策略尝试重新处理消息,之后将消息移至延迟队列(2)。

@Bean
public RabbitRetryTemplateCustomizer rabbitRetryTemplateCustomizer() {
    SimpleRetryPolicy retryPolicy = new SimpleRetryPolicy(3, Map.of(SkipIncrementRetryCounterException.class, true), true);

    ExponentialBackOffPolicy backOffPolicy = new ExponentialBackOffPolicy();
    backOffPolicy.setInitialInterval(Duration.ofSeconds(1).toMillis());
    backOffPolicy.setMaxInterval(Duration.ofSeconds(2).toMillis());
    backOffPolicy.setMultiplier(2);

    return (target, retryTemplate) -> {
        retryTemplate.setRetryPolicy(retryPolicy);
        retryTemplate.setBackOffPolicy(backOffPolicy);
    };
}
技术问询
  1. 当前实现方案是否存在优化空间?
  2. Spring框架是否有内置机制可实现该消息重试需求?
  3. 当前实现存在哪些潜在问题?

若有内容不清晰之处,请告知,我会进一步补充说明。


内容的提问来源于stack exchange,提问作者Piotr Olaszewski

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.21 15:27:57