详解RocketMQ的事务消息机制
·
RocketMQ 的事务事务消息机制是其核心特性之一,专门用于解决分布式系统中 “本地事务与远程操作一致性” 问题(即分布式事务问题)。它通过 “两阶段提交” 思想,确保本地业务操作与消息发送的最终一致性,避免出现 “本地事务成功但消息未发出” 或 “消息发出但本地事务失败” 的情况。
先帮助大家快速理解一下:
RocketMQ 的事务消息,简单说就是保证 “本地办事” 和 “发消息通知” 要么都成,要么都不成的机制。
比如你网购下单:
- 系统先扣库存(本地办事),再发消息让支付系统生成支付单(发通知)。
- 怕就怕扣了库存,发消息时服务器崩了 —— 库存没了,支付单没生成,订单废了。
事务消息怎么解决?分三步:
- 先 “占个坑”:发一条 “半消息” 给 MQ,这消息暂时藏着,别人看不到。
- 办正事:扣库存。成了就告诉 MQ “消息可以发出去了”;败了就告诉 MQ “把那条消息删了”。
- 怕失忆?回头查:如果扣完库存没来得及通知 MQ(比如服务器崩了),MQ 过一会儿会主动来问 “你刚才那事办没办成?”,系统查下库存记录,再告诉 MQ 结果。
这样就确保了:要么库存扣了、消息发了;要么库存没扣、消息也没发,两边一致,不添乱。
详::
核心场景:为什么需要事务消息?
举个典型例子:电商系统中 “下单减库存” 与 “发送支付通知” 的流程。
- 正确逻辑:先扣减库存(本地事务),再发送 “待支付” 消息通知用户付款。
- 问题:若扣减库存成功后,发送消息时服务器宕机,会导致用户收不到通知,订单异常。
- 事务消息作用:确保 “扣减库存” 和 “发送消息” 要么都成功,要么都失败(或回滚到一致状态)。
事务消息机制原理:两阶段提交 + 回查补偿
RocketMQ 事务消息的核心是将消息发送分为 “预处理” 和 “确认 / 回滚” 两个阶段,并通过 “消息回查” 机制解决本地事务结果未知的问题。整体流程如下:
1. 阶段一:发送半事务消息(Half Message)
- 半事务消息:一种特殊消息,发送到 RocketMQ 后,会被标记为 “暂不能投递” 状态,消费者无法消费。
- 作用:验证消息发送通道是否通畅,同时预留消息存储空间,为后续确认 / 回滚做准备。
流程:
- 生产者(业务系统)向 RocketMQ 发送 “半事务消息”(包含业务数据,如订单 ID、库存信息)。
- RocketMQ 收到消息后,存储到本地,并返回 “发送成功” 给生产者(此时消息处于 “半事务状态”,消费者不可见)。
2. 阶段二:执行本地事务 + 提交 / 回滚消息
生产者收到半事务消息发送成功的响应后,执行本地业务逻辑,并根据结果决定消息的最终状态。
流程:
- 生产者执行本地事务(如扣减库存、创建订单)。
- 根据本地事务结果,向 RocketMQ 发送 “提交” 或 “回滚” 指令:
- 若本地事务成功:发送 “提交” 指令,RocketMQ 将半事务消息标记为 “可投递”,消费者可正常消费。
- 若本地事务失败:发送 “回滚” 指令,RocketMQ 直接删除半事务消息,消费者不会收到。
3. 关键补偿:事务状态回查(解决未知状态)
若生产者在执行本地事务后,因网络中断、服务器宕机等原因,未向 RocketMQ 发送 “提交 / 回滚” 指令,RocketMQ 会主动发起 “事务回查”,确认本地事务的实际结果。
流程:
- RocketMQ 定期(默认 60 秒,可配置)扫描处于 “半事务状态” 的消息,若超过超时时间未收到确认指令,则触发回查。
- RocketMQ 调用生产者预设的 “回查接口”(由业务系统实现),查询本地事务的实际状态(成功 / 失败)。
- 生产者根据本地事务日志(需业务系统记录,如订单表中的 “事务状态” 字段),返回明确结果:
- 若本地事务成功:RocketMQ 提交消息,消费者可见。
- 若本地事务失败:RocketMQ 回滚消息,删除该消息。
- 若仍未知(如本地事务未执行完):RocketMQ 会在下次周期继续回查(默认最多回查 15 次,超过则丢弃消息)。
事务消息核心组件
- TransactionMQProducer:生产者类,专门用于发送事务消息,需配置 “本地事务执行器” 和 “回查处理器”。
- LocalTransactionExecuter:本地事务执行器,定义本地业务逻辑(如扣库存),返回 “提交 / 回滚 / 未知” 状态。
- TransactionCheckListener:事务回查监听器,定义回查逻辑(查询本地事务状态),用于 RocketMQ 回查时调用。
代码示例:事务消息发送与处理
以下是 Java 代码示例,展示事务消息的核心流程(基于 RocketMQ 4.x):
1. 定义生产者与事务处理器
import org.apache.rocketmq.client.producer.TransactionMQProducer;
import org.apache.rocketmq.client.producer.LocalTransactionExecuter;
import org.apache.rocketmq.client.producer.LocalTransactionState;
import org.apache.rocketmq.common.message.Message;
import org.apache.rocketmq.common.message.MessageExt;
public class TransactionProducer {
public static void main(String[] args) throws Exception {
// 1. 创建事务消息生产者
TransactionMQProducer producer = new TransactionMQProducer("transaction_producer_group");
producer.setNamesrvAddr("localhost:9876"); // NameServer 地址
// 2. 设置本地事务执行器(执行本地业务 + 返回状态)
producer.setTransactionExecuter(new LocalTransactionExecuter() {
@Override
public LocalTransactionState executeLocalTransactionBranch(Message msg, Object arg) {
try {
// 解析消息中的业务数据(如订单ID)
String orderId = new String(msg.getBody());
// 执行本地事务:扣减库存(示例逻辑)
boolean stockDeductSuccess = deductStock(orderId);
if (stockDeductSuccess) {
// 本地事务成功,返回“提交”
return LocalTransactionState.COMMIT_MESSAGE;
} else {
// 本地事务失败,返回“回滚”
return LocalTransactionState.ROLLBACK_MESSAGE;
}
} catch (Exception e) {
// 异常情况(如网络问题),返回“未知”,等待回查
return LocalTransactionState.UNKNOW;
}
}
});
// 3. 设置事务回查监听器(RocketMQ 回查时调用)
producer.setTransactionCheckListener((MessageExt msg) -> {
// 解析消息,查询本地事务状态(需从数据库中查询订单/库存状态)
String orderId = new String(msg.getBody());
boolean isTransactionSuccess = checkLocalTransactionStatus(orderId);
if (isTransactionSuccess) {
return LocalTransactionState.COMMIT_MESSAGE;
} else {
return LocalTransactionState.ROLLBACK_MESSAGE;
}
});
// 启动生产者
producer.start();
// 4. 发送半事务消息
Message msg = new Message(
"transaction_topic", // 主题
"order_tag", // 标签
"ORDER_12345".getBytes() // 消息体(订单ID)
);
producer.sendMessageInTransaction(msg, null); // 第二个参数为自定义业务参数
// 关闭生产者(实际生产环境不轻易关闭)
// producer.shutdown();
}
// 模拟扣减库存(本地事务)
private static boolean deductStock(String orderId) {
// 实际业务中操作数据库扣减库存,返回是否成功
System.out.println("扣减订单 " + orderId + " 的库存...");
return true; // 假设扣减成功
}
// 模拟查询本地事务状态(回查时调用)
private static boolean checkLocalTransactionStatus(String orderId) {
// 从数据库查询订单的事务状态(如订单表中“stock_deduct_status”字段)
System.out.println("回查订单 " + orderId + " 的库存扣减状态...");
return true; // 假设实际扣减成功
}
}
2. 消费者接收消息(常规消费逻辑)
import org.apache.rocketmq.client.consumer.DefaultMQPushConsumer;
import org.apache.rocketmq.client.consumer.listener.ConsumeConcurrentlyContext;
import org.apache.rocketmq.client.consumer.listener.ConsumeConcurrentlyStatus;
import org.apache.rocketmq.client.consumer.listener.MessageListenerConcurrently;
import org.apache.rocketmq.common.message.MessageExt;
import java.util.List;
public class TransactionConsumer {
public static void main(String[] args) throws Exception {
DefaultMQPushConsumer consumer = new DefaultMQPushConsumer("transaction_consumer_group");
consumer.setNamesrvAddr("localhost:9876");
consumer.subscribe("transaction_topic", "*"); // 订阅事务消息主题
// 注册消息监听器,接收并处理已提交的消息
consumer.registerMessageListener((MessageListenerConcurrently) (msgs, context) -> {
for (MessageExt msg : msgs) {
String orderId = new String(msg.getBody());
System.out.println("收到订单 " + orderId + " 的支付通知消息,开始通知用户...");
}
return ConsumeConcurrentlyStatus.CONSUME_SUCCESS; // 消费成功
});
consumer.start();
System.out.println("消费者启动成功,等待接收消息...");
}
}
关键特性与限制
- 最终一致性:事务消息不保证 “实时一致性”,但通过回查机制确保 “最终一致性”(最多回查 15 次,超时后消息被丢弃,需业务兜底)。
- 性能损耗:相比普通消息,多了 “半事务消息” 和 “回查” 步骤,性能略低,适合核心业务(如交易、支付)。
- 消息不可篡改:半事务消息一旦发送,内容不可修改,确保回查时数据一致。
- 依赖本地日志:回查机制依赖业务系统的本地事务日志(如订单状态表),需确保日志准确、可查询。
总结
RocketMQ 事务消息通过 “半事务消息预处理→本地事务执行→提交 / 回滚→回查补偿” 四步流程,解决了分布式系统中 “本地操作与消息发送” 的一致性问题。其核心是用 “两阶段提交”+“定时回查” 确保消息最终状态与本地事务一致,是分布式事务的轻量级、高可用解决方案。
更多推荐
所有评论(0)