基于DLQ的可重试AMQP消息处理方案优化问询
免责声明!
以下代码可正常运行,仅咨询实现的正确性及其他可行解决方案。
some-domain-name.q.some-service.incoming-messages- 该队列存储所有待处理的入站消息。
some-domain-name.q.some-service.incoming-messages.postponed- 该队列存储所有无法处理的延迟入站消息(如遇业务、连接或其他异常)。队列设置了TTL,到期后消息将被移回队列(1)。
some-domain-name.q.dead-letter- 主死信队列(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); }; }
- 当前实现方案是否存在优化空间?
- Spring框架是否有内置机制可实现该消息重试需求?
- 当前实现存在哪些潜在问题?
若有内容不清晰之处,请告知,我会进一步补充说明。
内容的提问来源于stack exchange,提问作者Piotr Olaszewski

