1. 首页
  2. rocketmq源码分析

023-二十三、RocketMQ源码分析之事务消息实现原理下篇-消息服务器Broker提交回滚事务实现原理

作者:唯有坚持不懈 | 出处:https://blog.csdn.net/prestigeding/article/details/78888290


RocketMQ事务消息阅读目录指引:
RocketMQ源码分析之从官方示例窥探RocketMQ事务消息实现基本思想
RocketMQ源码分析之RocketMQ事务消息实现原理上篇
RocketMQ源码分析之RocketMQ事务消息实现原理中篇—-事务消息状态回查
RocketMQ源码分析之事务消息实现原理下篇-消息服务器Broker提交回滚事务实现原理
RocketMQ事务消息实战


本文将重点分析 RocketMQ Broker 如何处理事务消息提交、回滚命令,其核心实现就是根据 commitlogOffset 找到消息,如果是提交动作,就恢复原消息的主题与队列,再次存入 commitlog 文件进而转到消息消费队列,供消费者消费,然后将原预处理消息存入一个新的主题RMQ_SYS_TRANS_OP_HALF_TOPIC,代表该消息已被处理;回滚消息与提交事务消息不同的是,提交事务消息会将消息恢复原主题与队列,再次存储在 commitlog 文件中。源码入口:
EndTransactionProcessor\#processRequest


OperationResult result = new OperationResult(); if (MessageSysFlag.TRANSACTION_COMMIT_TYPE == requestHeader.getCommitOrRollback()) { // @1 result = this.brokerController.getTransactionalMessageService().commitMessage(requestHeader); // @2 if (result.getResponseCode() == ResponseCode.SUCCESS) { // @3 RemotingCommand res = checkPrepareMessage(result.getPrepareMessage(), requestHeader); // @4 if (res.getCode() == ResponseCode.SUCCESS) { MessageExtBrokerInner msgInner = endMessageTransaction(result.getPrepareMessage()); // @5 msgInner.setSysFlag(MessageSysFlag.resetTransactionValue(msgInner.getSysFlag(), requestHeader.getCommitOrRollback())); msgInner.setQueueOffset(requestHeader.getTranStateTableOffset()); msgInner.setPreparedTransactionOffset(requestHeader.getCommitLogOffset()); msgInner.setStoreTimestamp(result.getPrepareMessage().getStoreTimestamp()); // @6 RemotingCommand sendResult = sendFinalMessage(msgInner); // @7 if (sendResult.getCode() == ResponseCode.SUCCESS) { this.brokerController.getTransactionalMessageService().deletePrepareMessage(result.getPrepareMessage()); // @8 } return sendResult; } return res; } }

代码@1:如果请求为提交事务,进入事务消息提交处理流程。
代码@2:提交消息,别被这名字误导了,该方法主要是根据 commitLogOffsetcommitlog 文件中查找消息返回 OperationResult 实例。
img_0914_01_1.png
private MessageExt prepareMessage :消息对象。
private int responseCode:查找结果。
private String responseRemark :错误提示。
代码@3:如果成功查找到消息,则继续处理,否则返回给客户端,消息未找到错误信息。
代码@4:验证消息必要字段。
● 验证消息的生产组与请求信息中的生产者组是否一致。
● 验证消息的队列偏移量(queueOffset)与请求信息中的偏移量是否一致。
● 验证消息的 commitLogOffset 与请求信息中的 CommitLogOffset 是否一致。

代码@5:调用 endMessageTransaction 方法,该方法主要的目的就是恢复事务消息的真实的主题、队列,并设置事务ID。

代码@6:设置消息的相关属性,这一步应该直接在 endMessageTransaction 中实现就好,统一恢复原消息的数量,特别关注的是取消了事务相关的系统标记。

代码@7:发送最终消息,其实现原理非常简单,调用MessageStore将消息存储在commitlog文件中,此时的消息,会被转发到原消息主题对应的消费队列,被消费者消费。

代码@8:删除预处理消息(prepare),其实是将消息存储在主题为:RMQ_SYS_TRANS_OP_HALF_TOPIC的主题中,代表这些消息已经被处理(提交或回滚)。
上述就是事务消息提交的流程,事务回滚类似,接下来大概分析一下事务消息回滚的流程。
EndTransactionProcessor\#processRequest


else if (MessageSysFlag.TRANSACTION_ROLLBACK_TYPE == requestHeader.getCommitOrRollback()) { result = this.brokerController.getTransactionalMessageService().rollbackMessage(requestHeader); // @1 if (result.getResponseCode() == ResponseCode.SUCCESS) { RemotingCommand res = checkPrepareMessage(result.getPrepareMessage(), requestHeader); if (res.getCode() == ResponseCode.SUCCESS) { this.brokerController.getTransactionalMessageService().deletePrepareMessage(result.getPrepareMessage()); // @2 } return res; } }

代码@1:回滚消息,其实内部就是根据 commitlogOffset 查找消息。

代码@2:将消息存储在RMQ_SYS_TRANS_OP_HALF_TOPIC中,代表该消息已被处理,与提交事务消息不同的是,提交事务消息会将消息恢复原主题与队列,再次存储在 commitlog 文件中。

事务消息在Broker服务端的提交回滚流程就介绍到这了。其核心实现就是根据 commitlogOffset 找到消息,如果是提交动作,就恢复原消息的主题与队列,再次存入 commitlog 文件进而转到消息消费队列,供消费者消费,然后将原预处理消息存入一个新的主题RMQ_SYS_TRANS_OP_HALF_TOPIC,代表该消息已被处理;回滚消息与提交事务消息不同的是,提交事务消息会将消息恢复原主题与队列,再次存储在 commitlog 文件中。

写完了如果写得有什么问题,希望读者能够给小编留言,也可以点击[此处扫下面二维码关注微信公众号](https://www.ycbbs.vip/?p=28 "此处扫下面二维码关注微信公众号")

看完两件小事

如果你觉得这篇文章对你挺有启发,我想请你帮我两个小忙:

  1. 关注我们的 GitHub 博客,让我们成为长期关系
  2. 把这篇文章分享给你的朋友 / 交流群,让更多的人看到,一起进步,一起成长!
  3. 关注公众号 「方志朋」,公众号后台回复「666」 免费领取我精心整理的进阶资源教程
  4. JS中文网,Javascriptc中文网是中国领先的新一代开发者社区和专业的技术媒体,一个帮助开发者成长的社区,是给开发者用的 Hacker News,技术文章由为你筛选出最优质的干货,其中包括:Android、iOS、前端、后端等方面的内容。目前已经覆盖和服务了超过 300 万开发者,你每天都可以在这里找到技术世界的头条内容。

    本文著作权归作者所有,如若转载,请注明出处

    转载请注明:文章转载自「 Java极客技术学习 」https://www.javajike.com

    标题:023-二十三、RocketMQ源码分析之事务消息实现原理下篇-消息服务器Broker提交回滚事务实现原理

    链接:https://www.javajike.com/article/1699.html

« 024-二十四、RocketMQ事务消息实战
022-二十二、RocketMQ源码分析之RocketMQ事务消息实现原理中篇—-事务消息状态回查»

相关推荐

QR code