最近项目中,出现了少量RabbitMQ消息重复消费的情况【大概率是由于网络波动或框架升级导致】,因此我们决定做下消息幂等处理。
问题分析:组内讨论的方案有大概两种:
一种是放Redis进行消息id保存和判重;
另外一种则是放MySQL 表里保存和判重;
由于MQ本身是异步处理的,因此就不会考虑使用缓存提升性能,另外我们打算将消息记录持久化存储,以便以后的问题排查分析。最终选择将消息id和消息内容持久化到MySQL表里。
SpringBoot中使用RabbitMQ进行消息发送和接收,但是查看阿里云控制台MessageId这一栏全是空。由于发送消息的时候未进行messageId的设置,导致控制台看不到MessageId,同时消费者也或取不到这个messageId。
代码示例: 发送示例:@Autowired
private RabbitTemplate rabbitTemplate;
// 消息发送,带MessageId
MessagePostProcessor messagePostProcessor = new MessagePostProcessor() {
@Override
public Message postProcessMessage(Message msg) throws AmqpException {
final MessageProperties messageProperties = msg.getMessageProperties();
// 消息id
messageProperties.setMessageId(UUID.randomUUID().toString());
return msg;
}
};
rabbitTemplate.convertAndSend(messageExchange(), messageRoutingKey(), message, messagePostProcessor);
接收示例:
//确认收到消息, false 当前消费者确认收到,true 所有消费者确认收到
channel.basicAck(deliveryTag, false);
// 增加消息幂等逻辑
final String messageId = message.getMessageProperties().getMessageId();
Boolean handleResult = handleIdempotentMessage(messageId,mqMessage);
if (!handleResult) {
log.debug("消息重复,messageId:{}", messageId);
return;
}
private Boolean handleIdempotentMessage(String messageId, Message message) {
// 这里逻辑是先进行MessageId查询;
LambdaQueryWrapper queryWrapper = new LambdaQueryWrapper<>();
queryWrapper.eq(MqData::getMessageId, messageId);
final Long messageCount = mqDataMapper.selectCount(queryWrapper);
if (messageCount > 0) {
// 消息存在重复,返回false
// 返回之前保存当前重复的消息
MqData mqData = new MqData();
mqData.setMessageId(messageId);
mqDataMapper.create(mqData);
return false;
}
return true;
}
参考:
https://help.aliyun.com/document_detail/184768.html
https://help.aliyun.com/document_detail/177412.htm
-------------欢迎各位留言交流学习,如有不正确的地方,请予以指正。【Q:981233589】



