本文介绍了使用Spring AMQP优化技术派RabbitMQ的实现方案。原方案采用连接池方式存在长连接问题和单线程处理瓶颈,优化后通过Spring AMQP内置的CachingConnectionFactory管理连接池和Channel缓存。主要改进包括:使用@RabbitListener实现消费监听、配置生产者确认和失败回调机制、引入Redis分布式锁防止重复消费。…
本文最初发布于 CSDN,现迁移至本站并做格式整理。内容保留原始观点与发布时间。
1. 原方案
原先技术派使用的是连接池的方式发送和消费消息,发送和消费结束后将连接返还给连接池
1.1. 原方案存在的问题
启用mq以后,在不点赞的时候只是10s输出一次processConsumerMsg cycle. 这是因为原先设计的就是十秒钟执行一次consumerMsg方法,AI说消费者应保持长连接状态,具体遇到的问题如下:
- 第一次执行这个方法的时候,consumer对象传入的是channel-1,在handleDelivery方法内部也应当由channel-1进行basicAck,但是第一次执行时候没有点赞,只是
channel.basicConsume(queueName, false, consumer) 告诉服务器:“有消息推给我”,然后channel-1在方法结束前close掉了
- 这之后点赞,会触发consumer的
handleDelivery方法,但是在里面进行channel.basicAck的时候需要的是已被关闭的channel-1,所以会报错: AlreadyClosedException
而且,while(true)事实上只有单线程处理消息,无法并行
1.2. 普通解决方案
RabbitmaAutoConfig初始化的时候会调用 processConsumerMsg()方法,去掉while使得初始化的时候有一个消费者在线:
1 2 3 4 5 6 7 8
| @Override public void processConsumerMsg() { log.info("Begin to processConsumerMsg."); consumerMsg(CommonConstants.EXCHANGE_NAME_DIRECT, CommonConstants.QUERE_NAME_PRAISE, CommonConstants.QUERE_KEY_PRAISE); }
|
consumerMsg()方法里面注释掉close的部分,这样能保证不再报AlreadyClosedException
2. 使用Spring AMQP优化
Spring AMQP内置了CachingConnectionFactory, 它帮我们维护了一个连接池和Channel 缓存,所以我们放弃旧方案,拥抱新技术🤗
主要优化有:
- Spring AMQP实现点赞通知,通过
RabbitListener注解实现消费监听
setConfirmCallback和setReturnsCallback实现 消息 -> 通过交换机 然后路由到 -> 相应的队列过程中生产者确认和失败回调的简单告警
- 利用Redis 分布式锁防止重复消费
一句话写到简历:基于 Spring AMQP 重构 RabbitMQ 集成方案,集成生产者确认+退回模式保证消息可靠送达,采用基于Redis TTL 的非阻塞式分布式锁实现消费端幂等与防抖设计, 配置死信队列实现异常消息隔离。
2.1. 修改依赖
paicoding-core/pom.xml
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15
|
<dependency> <groupId>org.springframework.boot</groupId> <artifactId>spring-boot-starter-amqp</artifactId> <version>2.7.1</version> </dependency>
|
2.2. 新增配置类
RabbitmqConfig配置了message的转换器、消费者并发数、点赞交换机和队列的配置及路由绑定、死信队列等,连接池应该是在connectionFactory里面。
之前UserFootServiceImpl里调用publishMsg的时候,最后一个参数传进来的是Json串,这里已经设置了jsonMessageConverter,以后直接传readUserFootDO就ok,缺点是接口类需要改一下参数类型
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60 61 62 63 64 65 66 67 68 69 70 71 72 73 74 75 76 77 78 79 80 81 82 83 84 85 86 87 88 89 90 91 92 93 94 95 96 97 98 99 100 101 102 103 104 105 106 107 108 109 110 111 112 113 114 115 116 117 118 119 120 121 122 123 124 125 126 127 128 129 130 131 132 133 134 135 136 137
| @Slf4j @Configuration public class RabbitmqConfig {
@Bean public MessageConverter jsonMessageConverter() { return new Jackson2JsonMessageConverter(); }
@Bean public RabbitTemplate rabbitTemplate(ConnectionFactory connectionFactory) { RabbitTemplate rabbitTemplate = new RabbitTemplate(connectionFactory); rabbitTemplate.setMessageConverter(jsonMessageConverter());
rabbitTemplate.setConfirmCallback((correlationData, ack, cause) -> { if (ack) { log.info("消息成功达到交换机:id={}", correlationData != null ? correlationData.getId() : null); } else { log.error("消息没成功到达交换机,id={}, cause={}", correlationData != null ? correlationData.getId() : null, cause); } });
rabbitTemplate.setReturnsCallback(returned -> { log.error("消息路由失败: exchange={}, routingKey={}, replyCode={}, replyText={}", returned.getExchange(), returned.getRoutingKey(), returned.getReplyCode(), returned.getReplyText()); }); return rabbitTemplate; }
@Bean public SimpleRabbitListenerContainerFactory rabbitListenerContainerFactory( ConnectionFactory connectionFactory) { SimpleRabbitListenerContainerFactory factory = new SimpleRabbitListenerContainerFactory(); factory.setConnectionFactory(connectionFactory); factory.setMessageConverter(jsonMessageConverter()); factory.setAcknowledgeMode(AcknowledgeMode.MANUAL); factory.setConcurrentConsumers(3); factory.setMaxConcurrentConsumers(10); factory.setPrefetchCount(10); return factory; }
@Bean public DirectExchange praiseExchange() { return new DirectExchange(CommonConstants.EXCHANGE_NAME_DIRECT, true, false); }
@Bean public Queue praiseQueue() { return QueueBuilder .durable(CommonConstants.QUERE_NAME_PRAISE) .withArgument("x-dead-letter-exchange", CommonConstants.EXCHANGE_NAME_DIRECT + ".dlx") .withArgument("x-dead-letter-routing-key", "dlx") .build(); }
@Bean public Binding praiseBinding(Queue praiseQueue, DirectExchange praiseExchange) { return BindingBuilder .bind(praiseQueue) .to(praiseExchange) .with(CommonConstants.QUERE_KEY_PRAISE); }
@Bean public DirectExchange deadLetterExchange() { return new DirectExchange(CommonConstants.EXCHANGE_NAME_DIRECT + ".dlx", true, false); }
@Bean public Queue deadLetterQueue() { return QueueBuilder .durable(CommonConstants.QUERE_NAME_PRAISE + ".dlq") .build(); }
@Bean public Binding deadLetterBinding(Queue deadLetterQueue, DirectExchange deadLetterExchange) { return BindingBuilder .bind(deadLetterQueue) .to(deadLetterExchange) .with("dlx"); } }
|
2.3. 启动类添加 @EnableRabbit 注解
2.4. 修改RabbitmqServiceImpl
类似RedisTemplate一样,这里使用的就是RabbitTemplate,分布式锁保证幂等性。
为什么上Redis TTL分布式锁?
AI老师: 在点赞通知场景中,我采用 Redis TTL 分布式锁主要是为了兼顾消费端幂等性与高频操作防抖。由于 RabbitMQ 存在网络波动导致的重复投递,且用户侧容易出现短时间内多次点击的情况,若直接穿透到数据库依靠唯一索引去重,不仅数据库压力大,还可能导致‘数据回滚了但通知却误发了’的业务不一致。通过 Redis 的 SETNX 原子命令配合 5 秒自动过期时间,我在缓存层构建了一个轻量级的快速失败(Fail-Fast)窗口,既能毫秒级拦截重复消息,保护了下游数据库,又实现了‘同一时间窗口内只触达一次通知’的友好交互体验。
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60 61 62 63 64 65 66 67 68 69 70 71 72 73 74 75 76 77 78 79 80 81 82 83 84 85 86 87 88 89 90 91 92 93 94 95 96 97 98 99 100 101 102 103 104 105 106 107 108 109 110
| @Slf4j @Service public class RabbitmqServiceImpl implements RabbitmqService {
@Autowired private RabbitTemplate rabbitTemplate;
@Autowired private RedisTemplate redisTemplate;
@Autowired private NotifyService notifyService;
public boolean enabled() { return true; }
public void publishMsg(String exchange, com.rabbitmq.client.BuiltinExchangeType exchangeType, String routingKey, UserFootDO message) { String messageId = message.getUserId() + "_" + message.getDocumentId(); try { rabbitTemplate.convertAndSend(exchange, routingKey, message, msg -> { msg.getMessageProperties().setMessageId(messageId); return msg; }); log.info("RabbitMQ publish success: exchange={}, routingKey={}, messageId={}", exchange, routingKey, messageId);
} catch (Exception e) { log.error("RabbitMQ publish failed: exchange={}, routingKey={}, messageId={}", exchange, routingKey, messageId, e); throw e; } }
@RabbitListener(queues = CommonConstants.QUERE_NAME_PRAISE) public void consumerMsg(UserFootDO message, Channel channel, @Header(AmqpHeaders.DELIVERY_TAG) long deliveryTag, Message amqpMessage) { log.info("RabbitMQ receive message: {}", message); String messageId = message.getUserId() + "_" + message.getDocumentId(); String lockKey = "rabbitmq:consume:lock:" + messageId;
try { Boolean lockAcquired = redisTemplate.opsForValue() .setIfAbsent(lockKey, "1", 5, TimeUnit.SECONDS);
if(Boolean.FALSE.equals(lockAcquired)){ log.warn("RabbitMQ message already processed, skip: messageId={}", messageId); channel.basicAck(deliveryTag, false); return; } notifyService.saveArticleNotify(message, NotifyTypeEnum.PRAISE);
channel.basicAck(deliveryTag, false); log.info("RabbitMQ consume success: userId={}, articleId={}", message.getUserId(), message.getDocumentId());
} catch (Exception e) { log.error("RabbitMQ consume failed: {}", message, e); try { redisTemplate.delete(lockKey); channel.basicNack(deliveryTag, false, false); } catch (IOException ioException) { log.error("RabbitMQ nack failed", ioException); } } }
}
|
2.5. 修改 application-rabbitmq.yml
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19
| rabbitmq: host: 127.0.0.1 port: 5672 username: guest password: guest virtual-host: / listener: simple: acknowledge-mode: manual prefetch: 10 concurrency: 3 max-concurrency: 10 retry: enabled: true initial-interval: 1000 max-attempts: 3 max-interval: 10000 template: mandatory: true
|
2.6. 删除或者注释掉之前的旧文件
原来的RabbitmqProperties、RabbitMqAutoConfig都加了注解,启动的时候都会运行,这里该注释的都注释掉
RabbitmqConnection、RabbitmqConnectionPool只能被主动调用,可以留着
另外,RabbitMQTest类也得注释一下,不然中间报错,AI说是因为 CommonConstants.QUERE_KEY_PRAISE重复了,确实是这样
1 2 3 4
| rabbitmqService.consumerMsg(CommonConstants.EXCHANGE_NAME_DIRECT, CommonConstants.QUERE_KEY_PRAISE, CommonConstants.QUERE_KEY_PRAISE);
|