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

如何在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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.25 06:31:43