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

如何封装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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.15 04:16:07