首页技术栈归档照片墙音乐日记随想收藏夹友链留言关于

MQ解决消息重复问题

写作时间:2026-07-08

MQ解决消息重复问题

消息生成者

@Service
public class AccountServiceImpl extends ServiceImpl<AccountMapper, Account> implements AccountService {

    @Autowired
    private RestTemplate restTemplate;

    @Autowired
    private DiscoveryClient discoveryClient;  //用来获取所有注册到nacos中的服务

    @Autowired
    private AccountFeignClient accountFeignClient;

    @Autowired
    private RabbitTemplate rabbitTemplate;

    @Autowired
    private RedisTemplate redisTemplate;

    @Override
    public Account login(String cardNo, String password) {
        return null;
    }

    /**
     * 转账
     * @param cardNoTo
     * @param cardNoFrom
     * @param tranMoney
     */
    public boolean transaction(String cardNoTo,String cardNoFrom,Double tranMoney) {
        //TODO 修改转出和转入账户的余额 account表

        //修改转出账户的余额
        QueryWrapper<Account> accountQueryWrapper =new QueryWrapper<>();
        accountQueryWrapper.eq("cardNo",cardNoFrom);
        Account accountFrom  = this.getOne(accountQueryWrapper);
        accountFrom.setBalance(accountFrom.getBalance()-tranMoney);
        this.updateById(accountFrom);

        TransactionRecord record = new TransactionRecord();
        record.setAmout(String.valueOf(tranMoney)); //转账金额
        record.setCardno(cardNoTo);  //转入
        record.setTransactiondate(new SimpleDateFormat("yyyy-MM-dd").format(new Date()));
        record.setTransactiontype("存款"); // 0取款 1存款

        //发送消息
        //这个参数是用来做消息的唯一标识
        //发布消息时使用,存储在消息的headers中
        // 消息发送,带MessageId
        String messageId = UUID.randomUUID().toString();
        MessagePostProcessor messagePostProcessor = new MessagePostProcessor() {
            @Override
            public Message postProcessMessage(Message msg) throws AmqpException {
                final MessageProperties messageProperties = msg.getMessageProperties();
                // 消息id
                messageProperties.setMessageId(messageId);
                return msg;
            }
        };
        rabbitTemplate.convertAndSend(RabbitmqConfig.BANK_EXCHANGE,
                RabbitmqConfig.BANK_ROUTINGKEY,
                record,
                messagePostProcessor);
        redisTemplate.opsForValue().set(messageId,false);
        return true;
    }
}

消息消费者

@Component
public class BankConsumer {
    @Autowired
    private TransactionRecordService recordService;
    @Autowired
    private RedisTemplate redisTemplate;
    /**
     * 监听mq消息
     */
    @RabbitHandler
    @RabbitListener(queuesToDeclare = {@Queue("bank_queue")}) //指定监听哪一个队列
    //@Header(AmqpHeaders.DELIVERY_TAG)long deliveryTag 从消息头中获取消息的投递顺序号
    public void receive(@Payload TransactionRecord record,
                        @Header(AmqpHeaders.DELIVERY_TAG)long deliveryTag,
                        @Header(AmqpHeaders.MESSAGE_ID) String messageId,
                        Channel channel) throws IOException {
        System.out.println("=======>"+record);
        System.out.println("correlationID=========>"+messageId);
        /*
         * 思考题:消费者已经收到了消息,但是在添加数据的时候发生了异常,转账记录没有生成,
         * 消息被删除了,导致数据不一致,
         * 这个问题怎么解决 :
         *  手动确认消息
        */
        //调用service中添加方法生成转账记录     解决消息重复消费
        //判断该消息是否已经消费
        String val = redisTemplate.opsForValue().get(messageId).toString();
        boolean flag = Boolean.parseBoolean(val);
        if(flag){
            System.out.println("该条消息已经被消费");
            return;
        }else{
            try{
                boolean result = recordService.addRecords(record);
                if(result){
                    //向redis中存入该条消息
                    redisTemplate.opsForValue().set(messageId,true);
                    //说明消息已经成功,发送确认成功应答给mq
                    channel.basicAck(deliveryTag,true);
                }else{
                    //说明消息没有成功消费,发送确认失败应答给mq
                    //第三个参数: 是否重新入列
                    //向redis中存入该条消息
                    channel.basicNack(deliveryTag,false,true);
                }
            }catch (Exception e){
                e.printStackTrace();
                //说明消息没有成功消费,发送确认失败应答给mq
                channel.basicNack(deliveryTag,false,true);
            }
        }

    }

}
avatar

yuanyourdomain

写代码,做研究,记录生活。

RECOMMENDED

MyBatis 动态 SQL

2026-07-08

Nginx 基础入门

2026-07-08

Maven 多模块与私服

2026-07-08