)
说实话事务消息是 RocketMQ 最复杂、也最令人惊叹的特性之一。它解决的是一道难题——在分布式系统中如何保证“发送消息”和“执行本地事务”要么一起成功要么一起失败今天这篇文章我们就来彻底搞懂事务消息。我会从它要解决的问题开始一步步深入到半消息、事务回查、状态机等核心机制。老规矩配合流程图一步一图。八、事务消息事务消息诞生的背景与要解决的问题我们先来看一个经典的分布式事务困境。你在电商系统里做一个“下单扣库存”的业务。订单服务和库存服务是两个独立的服务通过消息来解耦分布式事务困境订单服务创建订单发送扣库存消息消息发送成功但订单插入失败⚠️ 消息已经发出去了库存被扣了但订单不存在→ 数据不一致订单插入成功但消息发送失败⚠️ 订单创建了但库存没扣→ 超卖风险问题可以归结为一句话本地数据库操作订单插入和消息发送扣库存通知这两个操作不在同一个事务里。它们要么都成功要么都失败——但在分布式环境下这很难做到。常见的解决方案及其局限方案 怎么做 问题先发消息再执行本地事务 消息发出去后再插入订单 消息发成功但订单插入失败 → 库存被误扣先执行本地事务再发消息 插入订单后再发消息 订单插入成功但消息发送失败 → 超卖本地事务 本地消息表 在同一个数据库里存消息表定时扫描发送 依赖数据库轮询效率低有延迟RocketMQ 事务消息的解法RocketMQ 事务消息的核心思想是让消息的发送和本地事务的执行通过“半消息 事务回查”机制达成最终一致性。它不依赖本地消息表不需要轮询不需要额外的存储全部在 Broker 和 Producer 之间协调完成。事务消息的核心流程半消息Half Message什么是半消息Half Message半消息是事务消息的核心概念。它是一种暂不可见的消息——Producer 发送了这条消息但 Consumer 暂时还消费不到。半消息✅ 消息已写入 CommitLog有存储❌ 消息在 ConsumeQueue 中被标记为不可见⏳ 等待本地事务执行结果半消息 消息已存储但不可消费像个“待确认”的订单半消息的设计意图半消息的存在解决了“消息发送成功了但本地事务还没执行”这个中间状态的问题。如果半消息直接可见Consumer 就能消费到——但此时本地事务可能还没执行完数据状态不确定如果半消息不可见Consumer 就消费不到——等本地事务执行完再决定是提交还是回滚这样消息的可见性和本地事务的执行结果绑定在了一起。本地事务与半消息的协调机制事务消息的完整流程核心就是一个两阶段协调的过程成功失败COMMITROLLBACK业务发起事务消息阶段一发送半消息Producer 发送半消息到 BrokerBroker 存储半消息标记为不可消费Broker 返回半消息发送成功阶段二执行本地事务Producer 执行本地业务逻辑如插入订单本地事务结果阶段三提交/回滚Producer 向 Broker发送事务状态状态Broker 将半消息标记为可见Consumer 可以消费到消息Broker 删除半消息标记删除消息被丢弃不消费关键点半消息是第一步先发半消息再执行本地事务。如果半消息发送失败直接结束不用执行本地事务。本地事务决定最终状态本地事务成功 → COMMIT本地事务失败 → ROLLBACK。Broker 最终执行COMMIT 使消息可见ROLLBACK 使消息被丢弃。事务回查Check机制的实现原理上面我们讨论的是正常情况——本地事务执行完Producer 能正常告诉 Broker 是 COMMIT 还是 ROLLBACK。但如果本地事务执行完后Producer 宕机了或者网络断了呢Broker 怎么知道这条半消息最终应该被提交还是回滚答案就是事务回查Transaction Check机制。否是否是否是Broker 扫描半消息半消息状态是否为 UNKNOWN跳过Broker 向 Producer发起事务回查请求Producer 是否能正常响应标记为需要再次回查等待下一轮回查次数是否超限默认回滚防止消息堆积Producer 执行事务状态检查逻辑返回事务状态COMMIT / ROLLBACK / UNKNOWNBroker 根据状态执行提交或回滚回查机制的核心设计Broker 主动发起Broker 会定期扫描处于 UNKNOWN 状态的半消息默认每 60 秒一次Producer 被动响应Producer 需要实现 TransactionListener 的 checkLocalTransaction() 方法根据业务状态返回最终决定可配置重试回查失败会重试但如果超过配置的重试次数默认 15 次Broker 会默认回滚为什么需要回查回查解决了分布式系统中的不确定性问题——当 Producer 不可达时Broker 不能无限等待必须有最终决策。回查机制确保了半消息不会永远处于“悬而未决”的状态。事务消息的状态COMMIT、ROLLBACK、UNKNOWNRocketMQ 事务消息有三种状态由 Producer 的 TransactionListener 返回状态 含义 Broker 操作 Consumer 是否可见COMMIT 本地事务成功提交消息 将半消息标记为可见 ✅ 可见可消费ROLLBACK 本地事务失败回滚消息 删除半消息 ❌ 不可见被丢弃UNKNOWN 本地事务状态未知 暂不操作等待回查 ❌ 暂不可见等待决策public class OrderTransactionListener implements TransactionListener {Override public LocalTransactionState executeLocalTransaction(Message msg, Object arg) { try { // 执行本地事务插入订单 orderService.createOrder(arg); return LocalTransactionState.COMMIT_MESSAGE; } catch (Exception e) { return LocalTransactionState.ROLLBACK_MESSAGE; } } Override public LocalTransactionState checkLocalTransaction(MessageExt msg) { // 回查时根据业务数据判断事务是否成功 String orderId msg.getProperty(orderId); if (orderService.orderExists(orderId)) { return LocalTransactionState.COMMIT_MESSAGE; } return LocalTransactionState.ROLLBACK_MESSAGE; }}UNKNOWN 的使用场景本地事务执行过程中发生了异常无法确定是否成功本地事务执行时间很长需要异步确认依赖外部系统的状态需要等待外部回调UNKNOWN 状态触发的回查机制给了系统足够的容错空间。事务消息的完整时序图与执行流程把前面所有的内容串起来我们来看一张完整的时序图本地数据库ConsumerBrokerProducer业务应用本地数据库ConsumerBrokerProducer业务应用消息不可见不被消费此时不向 Broker 发送最终决定alt[多次回查仍返回UNKNOWN]loop[事务回查默认每 60 秒]alt[本地事务成功][本地事务失败][本地事务执行异常/超时]发送事务消息发送半消息状态PREPARED存储半消息不可消费记录事务状态半消息发送成功回调执行本地事务执行本地事务插入订单等7a. 事务提交成功8a. 返回 COMMIT9a. 提交事务消息10a. 标记消息可见11a. Consumer 消费到消息7b. 事务回滚8b. 返回 ROLLBACK9b. 回滚事务消息10b. 删除半消息8c. 返回 UNKNOWN9c. 发起回查请求10c. 检查本地事务状态11c. 查询业务数据12c. 返回查询结果13c. 返回 COMMIT/ROLLBACK/UNKNOWN14c. 提交最终事务状态15c. 超过回查次数限制默认回滚流程概括发送半消息 → 2. 执行本地事务 → 3. 根据结果决定状态 → 4. COMMIT/ROLLBACK/UNKNOWN → 5. UNKNOWN 触发回查 → 6. 最终确定状态。这个流程保证了要么本地事务成功且消息被消费要么本地事务失败且消息被丢弃。事务回查的时间参数配置事务回查涉及多个时间参数合理配置这些参数是事务消息稳定运行的关键参数 默认值 说明 配置位置transactionTimeout 6000ms 本地事务执行超时时间。超过这个时间Broker 认为事务状态未知触发回查 Producer 端transactionCheckInterval 60000ms Broker 回查半消息的间隔 Broker 端transactionCheckMax 15 次 最大回查次数。超过后默认回滚 Broker 端Broker 端配置transactionCheckInterval60000 # 回查间隔 60 秒transactionCheckMax15 # 最大回查 15 次transactionCheckTimeout6000 # 回查请求超时时间 6 秒Producer 端配置DefaultMQProducer producer new DefaultMQProducer();producer.setTransactionTimeout(6000); // 事务超时 6 秒配置建议transactionTimeout 要根据业务调整如果本地事务涉及多个数据库操作或远程调用需要设置足够的超时时间transactionCheckInterval 不宜过短频繁回查会增加 Broker 和 Producer 的负载transactionCheckMax 要结合业务容忍度核心业务可以设置更大的值如 30 次给系统更长的恢复时间事务消息的使用限制与注意事项事务消息虽强但并非万能。在使用时有几条重要的限制需要牢记限制一不支持延迟消息和批量消息事务消息不能和延迟消息、批量消息组合使用。事务消息本身就是“先存后发”的机制和延迟消息的“定时投递”在逻辑上冲突。限制二Consumer 必须处理重复消费事务消息虽然保证了“事务提交后消息可见”但在异常情况下如网络重传、RebalanceConsumer 仍然可能重复消费消息。业务方必须做好幂等处理。限制三事务消息的发送是同步的事务消息只支持同步发送不支持异步或单向发送。因为需要等待本地事务执行结果同步是最自然的模式。限制四事务消息不支持广播消费广播消费模式下每个 Consumer 独立消费事务消息的状态管理会变得复杂所以 RocketMQ 的事务消息不支持广播模式。限制五半消息对 Broker 有存储开销每条半消息在 COMMIT 或 ROLLBACK 之前都会占用 Broker 的存储空间CommitLog 和 ConsumeQueue。如果有大量半消息长时间 UNKNOWN会导致存储膨胀。使用限制❌ 不支持延迟消息❌ 不支持批量消息❌ 不支持异步/单向发送❌ 不支持广播消费⚠️ 半消息占用存储空间事务消息回查失败的处理策略回查不是万能的它也可能失败。我们需要理解各种失败场景及应对策略场景一Producer 节点宕机如果所有 Producer 节点都宕机了Broker 的回查请求无法送达。→ 策略配置 transactionCheckMax 和 transactionCheckInterval在 Producer 恢复前Broker 会持续尝试回查。如果超过最大次数仍无响应Broker 默认回滚。运维层面需要监控“超时未确认”的半消息数量。场景二回查请求超时Producer 收到了回查请求但因为网络或负载原因响应超时。→ 策略Broker 会重试回查下一次扫描时再次发起。Producer 端应确保 checkLocalTransaction() 方法的执行时间远小于 transactionCheckTimeout。场景三checkLocalTransaction 返回 UNKNOWNProducer 仍然无法确定事务状态比如依赖的数据库暂时不可用。→ 策略返回 UNKNOWN 后Broker 会在下一个回查周期再次询问。建议在业务层实现重试机制直到能确定状态为止。场景四超过最大回查次数15 次回查都失败了消息仍然 UNKNOWN。→ 策略Broker 默认回滚消息从 5.x 开始可以通过配置修改默认行为。运维需要监控这类情况手动介入处理。Producer 宕机响应超时返回 UNKNOWN超过最大次数回查失败失败类型等待 Producer 恢复Broker 持续重试超限后默认回滚优化网络或缩短 checkLocalTransaction 耗时优化业务检查逻辑确保能快速确定状态监控告警人工介入排查Done事务消息与分布式事务的对比TCC、Saga事务消息只是分布式事务的一种解决方案。我们来对比一下主流的几种方案TCCTry-Confirm-Cancel维度 TCC 事务消息核心思想 预留资源 → 确认/取消 半消息 → 提交/回滚业务侵入 高需要实现 Try/Confirm/Cancel 中只需实现本地事务 回查数据一致性 强一致性 最终一致性适用场景 跨服务调用、资源预留 消息驱动的异步场景性能 较低多次 RPC 较高一次半消息 本地事务Saga维度 Saga 事务消息核心思想 长事务拆分为多个本地事务 补偿 本地事务 消息最终一致性业务侵入 高需要实现正向操作 补偿操作 中数据一致性 最终一致性 最终一致性适用场景 长事务、跨多个服务的复杂流程 消息驱动的异步解耦场景对比总结选型建议分布式事务方案对比TCC强一致性业务侵入高Saga最终一致性需要补偿事务消息最终一致性通过回查自愈 事务消息消息驱动、异步解耦、对一致性要求不是强同步 TCC跨服务调用、需要强一致性、资源可预留 Saga长事务、复杂流程、需要有补偿逻辑选型建议事务消息适合消息驱动的异步场景对性能要求高能接受最终一致性TCC 适合跨服务调用的强一致性场景资源可以预留如扣减库存、冻结资金Saga 适合复杂的长事务流程每个步骤都有明确的补偿操作事务消息与普通消息的性能对比事务消息因为多了“半消息 本地事务 可能的回查”这些步骤性能自然不如普通消息。我们来量化对比一下对比维度 普通消息 事务消息 性能差距网络往返 1 次发送 2-3 次发送半消息 提交/回滚 可能的回查 约 2 倍延迟Broker 写入 1 次 CommitLog 写入 2 次 CommitLog 写入半消息 提交标记 约 2 倍存储 IO存储空间 消息体 消息体 事务状态记录 约额外 30%-50%CPU 开销 低 中需要事务状态管理和扫描 约 20%-30% 额外TPS 上限 ~10 万级 ~2-3 万级 约为普通消息的 1/3 - 1/5事务消息发送发送半消息写入半消息返回 ACK执行本地事务发送 COMMIT/ROLLBACK标记消息状态返回 ACK2 次网络往返 2 次写入本地事务 可能的回查普通消息发送发送请求写入 CommitLog返回 ACK1 次网络往返 1 次写入性能损耗的核心来源半消息额外写入每条事务消息至少多一次 CommitLog 写入事务状态维护Broker 需要维护和扫描半消息的状态网络往返增加从 1 次变成 2 次半消息 提交/回滚回查开销如果返回 UNKNOWN还有额外的回查负载 小贴士如果你的业务不需要强一致性保障普通消息 幂等消费可能是更好的选择。事务消息虽好但不要滥用——它只应该在“本地事务和消息发送必须保持一致”的场景下使用。事务消息的内部存储结构了解事务消息的存储细节能帮你更深入地理解它的工作原理事务消息存储结构回滚后标记为已删除不影响 CommitLog消息不可消费提交后在原 ConsumeQueue 位置标记为可见可能增加 Op 消息记录提交操作半消息存储CommitLog 中写入消息体 事务标记ConsumeQueue 中写入但标记为不可见TransactionStateService记录事务状态 UNKNOWN关键存储组件组件 作用RMQ_SYS_TRANS_HALF_TOPIC 半消息的存储 Topic所有事务消息的半消息都暂存于此TransactionStateService 管理事务状态的服务扫描半消息并触发回查Op 消息 记录 COMMIT/ROLLBACK 操作用于快速判断半消息的最终状态小结这篇文章我们彻底吃透了 RocketMQ 的事务消息机制通过 7 张流程图搞清楚了事务消息要解决的问题——分布式环境下本地事务和消息发送的一致性问题半消息的概念——暂不可见的中间状态解决了“先发消息还是先执行事务”的困境本地事务与半消息的协调——两阶段流程发送半消息 → 执行本地事务 → 提交/回滚事务回查机制——当 Producer 不可达时Broker 主动回查确保消息不会永远悬而未决三种事务状态——COMMIT、ROLLBACK、UNKNOWN 的含义和流转时间参数配置——transactionTimeout、transactionCheckInterval 等参数的作用使用限制——不支持延迟/批量/异步/广播等回查失败的处理策略——各种异常场景的应对方案与其他分布式事务方案的对比——TCC、Saga 的差异和选型建议