多实例部署Spring Boot MQTT监听器引发数据库重复数据的解决方案
解决MQTT多实例重复插入数据库的方案
方案1:使用MQTT共享订阅(推荐)
Amazon MQ支持MQTT 3.1.1的共享订阅机制,将订阅主题改为$share/[共享组名]/abc/#后,多个实例订阅同一共享组的主题,MQTT broker会将消息分发给组内单个实例,从根源避免所有实例都接收消息导致重复插入。
修改代码中的订阅主题:
// 将原订阅主题"abc/#"替换为共享订阅主题 subscribe("$share/chat-service-group/abc/#");
- 共享组名可自定义(如
chat-service-group),同一组内的实例会自动分摊消息消费 - 无需修改数据库或业务逻辑,仅调整订阅主题即可解决核心问题
方案2:数据库层面添加唯一性约束
通过在数据库表中添加唯一索引,强制保证数据唯一性,即使多个实例同时消费消息,数据库也会拒绝重复插入操作。
针对群聊消息(GroupChatMessage):
- 确认消息包含唯一标识字段(如
messageId,无此字段则需在消息体中新增) - 在实体类中标记唯一约束:
@Entity public class GroupChatMessage { // 其他字段... @Column(unique = true, nullable = false) private String messageId; // 消息唯一ID // getter/setter... }
- 执行数据库迁移生成唯一索引
- 代码中捕获唯一约束异常,忽略重复插入:
try { groupChatMessageRepository.save(groupChatMessage); } catch (DataIntegrityViolationException e) { logger.info("消息已存在,跳过插入: {}", groupChatMessage.getMessageId()); }
针对在线状态(UserOnlineStatusVO):
用户在线状态属于更新操作,JPA的save()方法会根据主键自动执行更新,但可给userId添加唯一索引,保证每个用户只有一条状态记录,避免冗余更新。
方案3:使用分布式锁控制消费
通过Redis等分布式锁工具,确保同一消息仅被一个实例处理。
- 引入Redis依赖(Spring Boot项目可使用
spring-boot-starter-data-redis) - 在消费逻辑中添加锁判断:
// 以消息唯一ID作为锁键 String lockKey = "mqtt:message:lock:" + groupChatMessage.getMessageId(); // 尝试获取锁,超时时间设为5秒,锁过期时间设为10秒 Boolean lockAcquired = redisTemplate.opsForValue().setIfAbsent(lockKey, "locked", 10, TimeUnit.SECONDS); if (Boolean.TRUE.equals(lockAcquired)) { try { // 执行插入数据库操作 groupChatMessageRepository.save(groupChatMessage); } finally { // 释放锁 redisTemplate.delete(lockKey); } } else { logger.info("其他实例正在处理该消息,跳过: {}", groupChatMessage.getMessageId()); }
- 锁的过期时间需大于消息处理的最大耗时,避免锁提前释放导致重复处理
- 该方案适合无法使用共享订阅的场景
方案4:优化MQTT客户端配置
- 合理设置QoS级别:QoS 0为最多投递一次,可减少重复投递概率;QoS 1/2需结合前面的方案避免重复插入
- 及时确认消息:确保客户端处理完消息后向broker发送确认,避免broker重复投递消息
内容的提问来源于stack exchange,提问作者Inder Bhushan Jha
相关产品推荐
相关产品推荐

