SpringCloud RocketMQ分布式事务
π大星的日常 人气:0消息队列实现分布式事务原理
首先让我们来看一下基于消息队列实现分布式事务的原理方案。
柔性事务
发送消息的服务有个OUTBOX数据表,在进行INSERT、UPDATE、DELETE 业务操作时也会给OUTBOX数据表INSERT一条消息记录,这样可以保证原子性,因为这是基于本地的ACID事务。
OUTBOX表充当临时消息队列,然后我们在引入一个消息中继(MessageRelay)的服务,由他从OUTBOX表中读取数据并发布消息到消息组件。
消息中继的实现可以很简单,只需要通过定时任务定期从OUTBOX表中拉取最新未发布的数据,获取到数据后将数据发送给消息组件,最后将完成发送的消息从OUTBOX表中删除即可,对于失败的消息可以根据业务规则进行重试。
RocketMQ的事务消息
RocketMQ本身已经支持事务消息,如果你们项目使用了RocketMQ,可以直接借助RocketMQ的事务消息实现分布式事务,我们先看一下RocketMQ事务消息的原理然后再借助RocketMQ来实现分布式事务。
RocketMQ采用了2PC的思想来实现了提交事务消息,同时增加一个补偿逻辑来处理二阶段超时或者失败的消息,如下图所示。
分布式事务
RocketMQ实现事务消息主要分为两个阶段:正常事务的发送及提交、事务信息的补偿流程
整体流程为:
正常事务发送与提交阶段
1、生产者发送一个半消息给MQServer(半消息是指消费者暂时不能消费的消息)
2、服务端响应消息写入结果,半消息发送成功
3、开始执行本地事务
4、根据本地事务的执行状态执行Commit或者Rollback操作
事务信息的补偿流程
1、如果MQServer长时间没收到本地事务的执行状态会向生产者发起一个确认回查的操作请求
2、生产者收到确认回查请求后,检查本地事务的执行状态
3、根据检查后的结果执行Commit或者Rollback操作
补偿阶段主要是用于解决生产者在发送Commit或者Rollback操作时发生超时或失败的情况。
RocketMQ事务流程关键
事务消息在一阶段对用户不可见
事务消息相对普通消息最大的特点就是一阶段发送的消息对用户是不可见的,也就是说消费者不能直接消费。这里RocketMQ的实现方法是原消息的主题与消息消费队列,然后把主题改成RMQ_SYS_TRANS_HALF_TOPIC
,这样由于消费者没有订阅这个主题,所以不会被消费。
如何处理第二阶段的失败消息?
在本地事务执行完成后会向MQServer发送Commit或Rollback操作,此时如果在发送消息的时候生产者出故障了,那么要保证这条消息最终被消费,MQServer会像服务端发送回查请求,确认本地事务的执行状态。
当然了rocketmq并不会无休止的的信息事务状态回查,默认回查15次,如果15次回查还是无法得知事务状态,RocketMQ默认回滚该消息。
消息状态 事务消息有三种状态:TransactionStatus.CommitTransaction:提交事务消息,消费者可以消费此消息
TransactionStatus.RollbackTransaction:回滚事务,它代表该消息将被删除,不允许被消费。
TransactionStatus.Unknown:中间状态,它代表需要检查消息队列来确定状态。
代码实现
业务需求:用户请求订单微服务order-service
接口删除订单(退货),删除订单时需要调用account-service
的方法给账户增加余额,一个典型的分布式事务问题。
基础配置
在Order-Service和Account-Service中引入Rocket消息组件
<dependency> <groupId>org.apache.rocketmq</groupId> <artifactId>rocketmq-spring-boot-starter</artifactId> </dependency>
在配置中心添加RocketMQ的相关配置
rocketmq:
name-server: xxx.xx.x.xx:9876
producer:
group: cloud-group
在OrderService服务中建立一张事务日志表rocketmq_transaction_log(作用稍后说)
发送半消息
Order-Service作为分布式事务开始的入口,在Service层我们给RocketMQ发送一条半消息
OrderController入口
/** * 根据订单号删除订单 * @param orderNo 订单编号 */ @PostMapping("/order/delete") public ResultData<String> delete(@RequestParam String orderNo){ log.info("delete order id is {}",orderNo); orderService.delete(orderNo); return ResultData.success("订单删除成功"); }
直接调用orderService的delete方法
OrderServiceImpl业务逻辑
@Override public void delete(String orderNo) { Order order = orderMapper.selectByNo(orderNo); //如果订单存在且状态为有效,进行业务处理 if (order != null && CloudConstant.VALID_STATUS.equals(order.getStatus())) { String transactionId = UUID.randomUUID().toString(); //如果可以删除订单则发送消息给rocketmq,让用户中心消费消息 rocketMQTemplate.sendMessageInTransaction("add-amount", MessageBuilder.withPayload( UserAddMoneyDTO.builder() .userCode(order.getAccountCode()) .amount(order.getAmount()) .build() ) .setHeader(RocketMQHeaders.TRANSACTION_ID, transactionId) .setHeader("order_id",order.getId()) .build() ,order ); } }
首先校验一下订单状态,然后使用rocketMQTemplate.sendMessageInTransaction()
发送事务消息。
sendMessageInTransaction方法有三个参数:
- destination:目的地(主题),这里发送给
add-amount
这个topic - message:发送给消费者的消息体,需要使用
MessageBuilder.withPayload()
来构建消息 - arg:参数
注意,这里我们生成了一个transactionId,并放在header中跟消息一起发送(这里实际也可以构造成一个对象,放在arg里进行发送),作用后面再讲!
消息封装实体UserAddMoneyDTO
@Data @NoArgsConstructor @AllArgsConstructor @Builder public class UserAddMoneyDTO { /** * 用户编码 */ private String userCode; /** * 金额 */ private BigDecimal amount; }
这个类生产者和消费者都需要用到,所以我直接丢到common包中,大家根据项目实际情况决定放哪。
执行本地事务与回查
MQServer收到半消息后会告诉生产者order-service确认收到半消息,这时候order-service需要执行本地事务,执行完本地事务后再告诉MQServer本地事务的执行状态,确认此消息究竟是Commit还是Rollback。
RocketMQ提供了RocketMQLocalTransactionListener
接口,本地事务监听器,这个接口类的实现如下:
第一个方法executeLocalTransaction
为执行本地事务;第二个方法checkLocalTransaction
为检查本地事务的执行状态,也就是回查动作。
我们需要实现RocketMQLocalTransactionListener
接口,在executeLocalTransaction
方法中执行本地事务,在执行checkLocalTransaction
回查方法时告诉RocketMQ到底该提交还是回滚。
这里大家思考一个问题,本地事务已经执行完成了,怎么去回查本地事务的执行结果呢?
答案如下:我们可以在执行本地事务的时候同时生成一条事务日志,让本地事务与日志事务在同一个方法中,同时添加@Transactional
注解,保证两个操作事务是一个原子操作。
这样如果事务日志表中有这个本地事务的信息,那就代表本地事务执行成功,需要Commit,相反如果没有对应的事务日志,则表示执行失败,需要Rollback。这就是为什么我们上面在OrderService中需要建立一张事务日志表的原因。
实现RocketMQLocalTransactionListener
接口,完成事务执行逻辑
/** * 监听事务消息 * @author javadaily */ @Slf4j @RocketMQTransactionListener @RequiredArgsConstructor(onConstructor = @__(@Autowired)) public class AddUserAmountListener implements RocketMQLocalTransactionListener { private final OrderService orderService; private final RocketMqTransactionLogMapper rocketMqTransactionLogMapper; /** * 执行本地事务 */ @Override public RocketMQLocalTransactionState executeLocalTransaction(Message message, Object arg) { log.info("执行本地事务"); MessageHeaders headers = message.getHeaders(); //获取事务ID String transactionId = (String) headers.get(RocketMQHeaders.TRANSACTION_ID); Integer orderId = Integer.valueOf((String)headers.get("order_id")); log.info("transactionId is {}, orderId is {}",transactionId,orderId); try{ //执行本地事务,并记录日志 orderService.changeStatuswithRocketMqLog(orderId, CloudConstant.INVALID_STATUS,transactionId); //执行成功,可以提交事务 return RocketMQLocalTransactionState.COMMIT; }catch (Exception e){ return RocketMQLocalTransactionState.ROLLBACK; } } /** * 本地事务的检查,检查本地事务是否成功 */ @Override public RocketMQLocalTransactionState checkLocalTransaction(Message message) { MessageHeaders headers = message.getHeaders(); //获取事务ID String transactionId = (String) headers.get(RocketMQHeaders.TRANSACTION_ID); log.info("检查本地事务,事务ID:{}",transactionId); //根据事务id从日志表检索 QueryWrapper<RocketmqTransactionLog> queryWrapper = new QueryWrapper<>(); queryWrapper.eq("transaction_id",transactionId); RocketmqTransactionLog rocketmqTransactionLog = rocketMqTransactionLogMapper.selectOne(queryWrapper); if(null != rocketmqTransactionLog){ return RocketMQLocalTransactionState.COMMIT; } return RocketMQLocalTransactionState.ROLLBACK; } }
本地事务执行逻辑
@Transactional(rollbackFor = RuntimeException.class) @Override public void changeStatuswithRocketMqLog(Integer id,String status,String transactionId){ orderMapper.changeStatus(id,status); rocketMqTransactionLogMapper.insert( RocketmqTransactionLog.builder() .transactionId(transactionId) .log("执行删除订单操作") .build() ); }
修改订单状态为删除状态,同时往事务日志表中插入一条事务日志,用@Transactional注解保证事务。
Account-Service消费消息
监听消息并处理给用户增加余额逻辑
@Slf4j @Service @RocketMQMessageListener(topic = "add-amount",consumerGroup = "cloud-group") @RequiredArgsConstructor(onConstructor = @__(@Autowired) ) public class AddUserAmountListener implements RocketMQListener<UserAddMoneyDTO> { private final AccountMapper accountMapper; /** * 收到消息的业务逻辑 */ @Override public void onMessage(UserAddMoneyDTO userAddMoneyDTO) { log.info("received message: {}",userAddMoneyDTO); accountMapper.increaseAmount(userAddMoneyDTO.getUserCode(),userAddMoneyDTO.getAmount()); log.info("add money success"); } }
测试
测试数据
订单表
用户表
事务日志表
如果事务消息成功消费最终用户表中jianzh5这个用户的amount应该变成300(100+200)
测试准备
我们在执行本地事务成功并需要通知消息队列提交事务处打个断点,然后在执行到此处时手动模拟异常
模拟异常
在准备提交事务时我们通过命令taskkill /pid 10116 -t -f
命令强制杀掉OrderService进程。(先通过jps获取OrderService进程ID)
重启服务器,检查是否会执行回查方法
重启OrderService程序会自动执行回查方法,结合事务日志表判断是否提交事务。
运行后的结果
小结
我们介绍了使用消息队列实现柔性事务的方案,重点剖析了RocketMQ事务消息的原理,并通过Demo案例实现了分布式事务(柔性事务)。
加载全部内容