MQ解决消息重复问题
消息生成者
@Service
public class AccountServiceImpl extends ServiceImpl<AccountMapper, Account> implements AccountService {
@Autowired
private RestTemplate restTemplate;
@Autowired
private DiscoveryClient discoveryClient;
@Autowired
private AccountFeignClient accountFeignClient;
@Autowired
private RabbitTemplate rabbitTemplate;
@Autowired
private RedisTemplate redisTemplate;
@Override
public Account login(String cardNo, String password) {
return null;
}
public boolean transaction(String cardNoTo,String cardNoFrom,Double tranMoney) {
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("存款");
String messageId = UUID.randomUUID().toString();
MessagePostProcessor messagePostProcessor = new MessagePostProcessor() {
@Override
public Message postProcessMessage(Message msg) throws AmqpException {
final MessageProperties messageProperties = msg.getMessageProperties();
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;
@RabbitHandler
@RabbitListener(queuesToDeclare = {@Queue("bank_queue")})
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);
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){
redisTemplate.opsForValue().set(messageId,true);
channel.basicAck(deliveryTag,true);
}else{
channel.basicNack(deliveryTag,false,true);
}
}catch (Exception e){
e.printStackTrace();
channel.basicNack(deliveryTag,false,true);
}
}
}
}