如何封装Spring AMQP的@RabbitListener实现消息重发与异常处理
实现自定义
@RabbitResenderListener注解方案 这个需求很常见,咱们可以通过Spring AOP结合Spring AMQP的扩展能力来实现,把原来的重试和异常处理逻辑完全封装到自定义注解里,业务代码只需要专注自己的核心逻辑。下面是具体的实现步骤:
1. 定义自定义注解@RabbitResenderListener
首先创建自己的注解,复用@RabbitListener的核心配置(比如队列名),同时添加重试相关的自定义参数,让配置更灵活:
import org.springframework.amqp.rabbit.annotation.RabbitListener; import java.lang.annotation.*; @Target(ElementType.METHOD) @Retention(RetentionPolicy.RUNTIME) @Documented // 继承@RabbitListener元数据,确保Spring AMQP能识别队列配置 @RabbitListener public @interface RabbitResenderListener { // 复用@RabbitListener的队列配置属性 String[] queues() default {}; // 自定义最大重试次数,默认3次 int maxRetryCount() default 3; // 自定义重试间隔基数(指数退避:base^n秒),默认2 int baseRetryInterval() default 2; }
2. 实现AOP切面封装重试与异常处理逻辑
创建切面类,拦截所有被@RabbitResenderListener标注的方法,把原来的重试计数、手动确认、异常分支逻辑都封装在这里:
import com.rabbitmq.client.Channel; import org.aopalliance.intercept.MethodInterceptor; import org.aopalliance.intercept.MethodInvocation; import org.springframework.amqp.core.Message; import org.springframework.amqp.core.MessageProperties; import org.springframework.amqp.support.AmqpHeaders; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.data.redis.core.RedisTemplate; import org.springframework.messaging.handler.annotation.Header; import org.springframework.stereotype.Component; import org.springframework.aop.support.AopUtils; @Component public class RabbitResenderAspect implements MethodInterceptor { @Autowired private RedisTemplate<String, Object> redisTemplate; // 重试计数存储的Redis键值,可抽成配置类 private static final String MQ_RETRY_COUNT_KEY = "mq:consumer:retry:count"; @Override public Object invoke(MethodInvocation invocation) throws Throwable { // 1. 解析方法参数,获取Message、Channel、DeliveryTag核心对象 Message message = null; Channel channel = null; Long deliveryTag = null; // 获取当前方法的自定义注解配置 RabbitResenderListener annotation = AopUtils.getTargetClass(invocation.getThis()) .getMethod(invocation.getMethod().getName(), invocation.getMethod().getParameterTypes()) .getAnnotation(RabbitResenderListener.class); // 遍历参数匹配所需对象 for (int i = 0; i < invocation.getArguments().length; i++) { Object arg = invocation.getArguments()[i]; if (arg instanceof Message) { message = (Message) arg; } else if (arg instanceof Channel) { channel = (Channel) arg; } else if (arg instanceof Long && invocation.getMethod().getParameterAnnotations()[i] != null) { // 匹配@Header注解标注的DeliveryTag参数 for (java.lang.annotation.Annotation ann : invocation.getMethod().getParameterAnnotations()[i]) { if (ann instanceof Header && ((Header) ann).value().equals(AmqpHeaders.DELIVERY_TAG)) { deliveryTag = (Long) arg; } } } } // 参数校验,确保核心对象存在 if (message == null || channel == null || deliveryTag == null) { throw new IllegalArgumentException("方法参数必须包含Message、Channel和@Header(AmqpHeaders.DELIVERY_TAG)的Long类型参数"); } MessageProperties properties = message.getMessageProperties(); String messageId = properties.getMessageId(); int maxRetryCount = annotation.maxRetryCount(); int baseInterval = annotation.baseRetryInterval(); // 2. 记录当前重试次数 Long retryCount = redisTemplate.opsForHash().increment(MQ_RETRY_COUNT_KEY, messageId, 1); try { // 3. 执行业务方法 Object result = invocation.proceed(); // 4. 消息处理成功,手动确认并清除重试计数 channel.basicAck(deliveryTag, false); redisTemplate.opsForHash().delete(MQ_RETRY_COUNT_KEY, messageId); return result; } catch (Exception e) { // 5. 异常分支处理 if (retryCount >= maxRetryCount) { // 超过最大重试次数,拒绝消息(不再重新入队) channel.basicReject(deliveryTag, false); redisTemplate.opsForHash().delete(MQ_RETRY_COUNT_KEY, messageId); } else { // 指数退避延迟后,将消息重新入队 long delay = (long) (Math.pow(baseInterval, retryCount) * 1000); Thread.sleep(delay); channel.basicNack(deliveryTag, false, true); } throw e; } } }
3. 配置AOP代理让切面生效
通过Bean后置处理器,为带有@RabbitResenderListener注解的Bean添加切面代理,确保逻辑能被拦截:
import org.springframework.aop.framework.ProxyFactory; import org.springframework.aop.support.AopUtils; import org.springframework.beans.BeansException; import org.springframework.beans.factory.config.BeanPostProcessor; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; import org.springframework.aop.framework.autoproxy.DefaultAdvisorAutoProxyCreator; @Configuration public class RabbitResenderConfig { @Autowired private RabbitResenderAspect rabbitResenderAspect; // 开启自动代理,支持类代理 @Bean public DefaultAdvisorAutoProxyCreator defaultAdvisorAutoProxyCreator() { DefaultAdvisorAutoProxyCreator creator = new DefaultAdvisorAutoProxyCreator(); creator.setProxyTargetClass(true); return creator; } // 拦截RabbitListener Bean的初始化,为其添加切面 @Bean public BeanPostProcessor rabbitListenerBeanPostProcessor() { return new BeanPostProcessor() { @Override public Object postProcessAfterInitialization(Object bean, String beanName) throws BeansException { if (AopUtils.isAopProxy(bean)) { return bean; } // 检查当前Bean是否有被@RabbitResenderListener标注的方法 boolean hasResenderAnnotation = false; for (java.lang.reflect.Method method : bean.getClass().getDeclaredMethods()) { if (method.isAnnotationPresent(RabbitResenderListener.class)) { hasResenderAnnotation = true; break; } } if (hasResenderAnnotation) { ProxyFactory proxyFactory = new ProxyFactory(bean); proxyFactory.addAdvice(rabbitResenderAspect); return proxyFactory.getProxy(); } return bean; } }; } }
4. 业务代码简洁使用示例
现在业务类只需要标注自定义注解,编写核心业务逻辑即可,重试、确认、异常处理都由切面自动完成:
import com.rabbitmq.client.Channel; import org.springframework.amqp.core.Message; import org.springframework.amqp.support.AmqpHeaders; import org.springframework.messaging.handler.annotation.Header; import org.springframework.stereotype.Component; @Component public class BizMessageConsumer { @RabbitResenderListener(queues = "so38728668", maxRetryCount = 5) public void handleMessage(Message message, Channel channel, @Header(AmqpHeaders.DELIVERY_TAG) long deliveryTag) { // 这里只需要处理核心业务逻辑,无需关心重试和消息确认 System.out.println("处理消息内容:" + new String(message.getBody())); // 模拟业务异常(测试重试逻辑) // throw new RuntimeException("业务处理失败,触发重试"); } }
关键说明
- 切面完全封装了原有的重试计数、手动确认、指数退避逻辑,业务代码无需继承抽象类,侵入性极低。
- 自定义注解继承了
@RabbitListener,Spring AMQP可以正常识别队列配置,无需额外绑定配置。 - 重试次数通过Redis哈希表持久化,避免服务重启后重试计数丢失(如果不需要持久化,也可以改用本地缓存)。
内容的提问来源于stack exchange,提问作者vvsueprman
相关产品推荐
相关产品推荐

