如何在Spring AMQP中向RabbitMQ队列发送唯一消息?
解决Spring AMQP+RabbitMQ重复消息导致数据库主键冲突的方案
这个场景我太熟悉了——重复消息+多节点并发消费+主键约束,绝对是分布式消息架构里的典型坑。给你几个实操性拉满的解决方案,按落地难度和可靠性排序,你可以根据业务场景组合使用:
1. 从源头掐断:生产者端去重
既然重复消息是因为服务被多次调用导致的,那咱们先在生产者环节把重复请求拦下来:
- 核心思路:给每条待发送的消息生成唯一业务ID(比如对数据内容做MD5哈希,或者用「业务标识+UUID」的组合),发送前先把这个ID存入Redis(用原子操作保证唯一性),只有写入成功才发消息。
- 代码示例:
// 生成唯一业务ID,这里用数据ID+时间戳做示例,你可以换成更贴合业务的规则 String bizMsgId = data.getId() + "_" + System.currentTimeMillis(); // 用Redis的setIfAbsent做原子判断,设置24小时过期避免存储膨胀 Boolean isFirstSend = stringRedisTemplate.opsForValue() .setIfAbsent(bizMsgId, "sent", 24, TimeUnit.HOURS); if (Boolean.TRUE.equals(isFirstSend)) { // 首次发送,正常投递消息 rabbitTemplate.convertAndSend("your-exchange", "your-routing-key", data); } else { // 重复请求,直接返回不发消息 log.info("重复数据,跳过发送:{}", data); }
2. RabbitMQ消费者端:幂等性消费
就算生产者漏了,消费者也要自己做好防护,核心是记录已处理的消息ID,避免重复执行:
- 方案一:利用消息头的唯一ID+Redis做幂等校验
@RabbitListener(queues = "your-queue") public void handleMessage(Message message, Channel channel) throws IOException { // 从消息头获取自定义的业务ID(生产者发送时要记得设置) String bizMsgId = message.getMessageProperties().getHeaders().get("biz-msg-id").toString(); Long deliveryTag = message.getMessageProperties().getDeliveryTag(); // 原子判断是否已处理 Boolean isProcessed = stringRedisTemplate.opsForValue() .setIfAbsent(bizMsgId, "processed", 24, TimeUnit.HOURS); if (Boolean.TRUE.equals(isProcessed)) { try { // 执行数据库插入逻辑 yourDataMapper.insert(data); // 处理成功,手动确认消息 channel.basicAck(deliveryTag, false); } catch (DuplicateKeyException e) { // 兜底:就算Redis校验漏了,数据库抛异常也标记为已处理 stringRedisTemplate.opsForValue().set(bizMsgId, "processed", 24, TimeUnit.HOURS); channel.basicAck(deliveryTag, false); log.warn("数据已存在,跳过处理:{}", data); } } else { // 已处理过,直接确认消息 channel.basicAck(deliveryTag, false); log.info("消息已处理,跳过:{}", bizMsgId); } } - 方案二:开启手动确认模式,避免消息重复投递
调整Spring AMQP的容器配置,确保只有处理成功才确认消息,失败则根据情况重试或进入死信队列:@Bean public SimpleRabbitListenerContainerFactory rabbitListenerContainerFactory(ConnectionFactory connectionFactory) { SimpleRabbitListenerContainerFactory factory = new SimpleRabbitListenerContainerFactory(); factory.setConnectionFactory(connectionFactory); factory.setConcurrentConsumers(4); // 按节点数调整,避免过度并发 factory.setAcknowledgeMode(AcknowledgeMode.MANUAL); // 手动确认 return factory; }
3. 数据库层面:最后一道兜底防线
前面的环节都有概率漏网,数据库必须做最后一层保障:
- 方案一:使用UPSERT语法(MySQL/PostgreSQL都支持)
把插入操作转为「不存在则插入,存在则更新」,直接避免主键冲突:-- MySQL示例 INSERT INTO your_table (id, content, create_time) VALUES (#{id}, #{content}, NOW()) ON DUPLICATE KEY UPDATE content = #{content}, update_time = NOW(); - 方案二:事务内的查询+插入
如果不能用UPSERT,就把「查询是否存在」和「插入」放在同一个事务里,用数据库事务的隔离性避免并发冲突:
注意:要把事务隔离级别设为@Transactional(rollbackFor = Exception.class) public void insertIfNotExists(Data data) { Data existing = yourDataMapper.selectById(data.getId()); if (existing == null) { yourDataMapper.insert(data); } }REPEATABLE READ,避免幻读问题。
总结建议
推荐组合使用「生产者去重+消费者Redis幂等+数据库UPSERT」,从源头到末端形成完整的幂等防护链。毕竟分布式系统里,重复消息是无法完全避免的,咱们只能通过层层设防来保证业务的一致性。
内容的提问来源于stack exchange,提问作者user1159553
相关产品推荐
相关产品推荐

