彻底搞懂 RocketMQ 事务消息:从两阶段提交原理到分布式实战(超详细)
🚀 彻底搞懂 RocketMQ 事务消息:从两阶段提交原理到分布式实战(超详细)
在分布式系统开发中,我们经常遇到这样的难题:“本地数据库操作(如生成订单)”与“发送 MQ 消息(如通知积分系统发积分)”如何保证同时成功或同时失败?
如果订单生成成功,但网络抖动导致 MQ 消息没发出去,用户就会投诉没拿到积分;反之,如果订单因库存不足失败了,MQ 消息却发出去了,积分系统就会凭空给用户送积分。
为了解决这个“分布式事务”难题,RocketMQ 引入了极其强大的事务消息(Transaction Message)机制。本文将用最接地气的“生活类比”,带你从零彻底吃透它的原理与核心代码。
一、 核心原理:两阶段提交与回查机制
很多初学者一听到“事务”和“两阶段提交”就头大。别慌,我们用一个生活场景来类比:你在网上买东西,付了钱,商家才给你发货。 RocketMQ 事务消息的精髓可以总结为九个字:“先斩后奏” →\rightarrow→ “做自家的事” →\rightarrow→ “补发通知”。

1. 正常主线流程(一阶段 + 二阶段)
- 步骤 1. Send: Half Msg(发送半消息)
大白话:你准备给朋友转账,先给银行打个招呼:“我准备转 100 块钱,你先帮我把这笔交易挂起,但千万别把钱真的放进对方账户。”
技术解析:生产者(Producer)先向 RocketMQ Broker 发送一条“半消息”。此时消息已经到了 Broker,但被存放在一个特殊队列中,消费者(Subscriber)不可见。
- 步骤 2. Half Msg Send OK(半消息成功响应)
大白话:银行回复你:“好的,单子我先拿在手里了,对方现在还看不到这笔钱。你可以去干你的正事了。”
- 步骤 3. Execute Local Transaction(执行本地事务)
大白话:你收到银行确认后,赶紧在自己的小账本上写下:“扣除本地余额 100 元”。
技术解析:生产者开始执行本地数据库操作(如向 t_order 表插入订单)。
- 步骤 4. Commit or Rollback(最终提交或回滚)
大白话 & 技术解析:这是最关键的抉择点,根据本地事务的结果做出选择:
Commit(提交):本地数据库成功了!告诉 Broker 把刚才的“半消息”变成“正式消息”,此时消费者才能看到并领走这条消息。
Rollback(回滚):本地数据库失败了(如余额/库存不足)!告诉 Broker 把半消息直接删掉(Delete),消费者永远不会知道这件事。
2. 兜底补偿流程:防止失联的“回查机制”
如果在执行完步骤 3(本地扣完钱了),正准备告诉 Broker 去 Commit 的时候,生产者的电脑突然断电或断网了怎么办? 这条半消息难道一直悬挂在 Broker 里吗?
RocketMQ 给出了完美的防丢兜底方案 —— 回查机制(Check back):
- 步骤 5. Check back(服务器定时回查)
大白话:银行系统等了半天没收到你的最终确认,主动打电话给你:“刚刚看你要转账,后来你那边没信号了,你本地到底扣钱成功了没有?”
技术解析:Broker 发现某条“半消息”超时未决,会主动发起网络请求,回查该生产者的本地状态。
- 步骤 6. Check the state(检查本地事务状态)
大白话:你接到电话,翻开自己的小账本,确认刚刚那笔账有没有写成功。
技术解析:生产者程序收到回查后,查询本地数据库,看订单记录是否存在。
- 步骤 7. Re-Commit or Rollback(根据回查结果再次确认)
大白话:你告诉银行:“我翻了账本,本地确实扣钱成功了,请把钱放过去吧!”
技术解析:生产者根据查询结果,再次向 Broker 发送 Commit 或 Rollback 指令。
二、 实战代码演练:电商下单并发放积分
下面我们用 Java 语言实现这一全套逻辑。我们将模拟两个订单:
- 订单 1002(偶数):本地事务执行成功 →\rightarrow→ 期望最终 Commit(消费者能收到积分消息)。
- 订单 1003(奇数):本地事务执行失败 →\rightarrow→ 期望最终 Rollback(消费者收不到消息)。
1. 编写事务监听器(TransactionListener)
核心逻辑都在这里:executeLocalTransaction 负责下单,checkLocalTransaction 负责应对 Broker 的“肉体回查”。
import org.apache.rocketmq.client.producer.LocalTransactionState;
import org.apache.rocketmq.client.producer.TransactionListener;
import org.apache.rocketmq.common.message.Message;
import org.apache.rocketmq.common.message.MessageExt;
import java.util.concurrent.ConcurrentHashMap;
public class OrderTransactionListener implements TransactionListener {
// 模拟一个本地数据库。Key: 订单ID, Value: 1代表成功, 2代表失败
private ConcurrentHashMap<String, Integer> localDatabase = new ConcurrentHashMap<>();
/**
* 【核心方法 A】执行本地事务(对应图中 步骤3)
* 当半消息(Half Msg)发送成功后,RocketMQ 会自动调用该方法
*/
@Override
public LocalTransactionState executeLocalTransaction(Message msg, Object arg) {
String orderId = msg.getKeys(); // 拿到订单ID作为唯一业务标识
System.out.println("【本地事务】收到半消息成功响应,开始在本地数据库下单,订单ID: " + orderId);
try {
// =======================================================
// 模拟真实的业务:这里写你的真实数据库操作(如:insert into t_order...)
// =======================================================
System.out.println("【本地事务】正在向数据库写入订单信息...");
Thread.sleep(500); // 模拟耗时
// 模拟业务规则:假设偶数订单成功,奇数订单库存不足失败
if (Integer.parseInt(orderId) % 2 == 0) {
localDatabase.put(orderId, 1); // 1 代表成功保存到数据库
System.out.println("【本地事务】数据库写入成功!准备向 Broker 发送 Commit");
// 对应 步骤4:告诉服务器这笔交易成了,消息可以投递给消费者了
return LocalTransactionState.COMMIT_MESSAGE;
} else {
localDatabase.put(orderId, 2); // 2 代表失败
System.out.println("【本地事务】库存不足,数据库写入失败!准备向 Broker 发送 Rollback");
// 对应 步骤4:告诉服务器这笔交易废了,直接把这条半消息删掉
return LocalTransactionState.ROLLBACK_MESSAGE;
}
} catch (Exception e) {
System.out.println("【本地事务】发生未知异常!暂不回应 Broker,等待回查...");
// 如果程序中间断电、崩溃或返回 UNKNOW,Broker 会在超时后自动触发【回查机制】
return LocalTransactionState.UNKNOW;
}
}
/**
* 【核心方法 B】本地事务回查(对应图中 步骤5、6、7)
* 如果上面的方法由于网络断开没有成功返回,或者返回了 UNKNOW,Broker 会定期自动调用此方法
*/
@Override
public LocalTransactionState checkLocalTransaction(MessageExt msg) {
String orderId = msg.getKeys();
System.out.println("【事务回查】Broker 发现状态不明,主动来回查了!正在查询订单ID: " + orderId);
// 翻看本地账本(查询本地数据库状态)
Integer status = localDatabase.get(orderId);
if (status != null && status == 1) {
System.out.println("【事务回查】检查账本确认:该订单之前确实成功了。回应:Commit!");
return LocalTransactionState.COMMIT_MESSAGE;
} else if (status != null && status == 2) {
System.out.println("【事务回查】检查账本确认:该订单之前确实失败了。回应:Rollback!");
return LocalTransactionState.ROLLBACK_MESSAGE;
}
System.out.println("【事务回查】账本里依然找不到这个订单!回应:继续保持 UNKNOW");
return LocalTransactionState.UNKNOW;
}
}
2. 编写事务消息生产者(TransactionMQProducer)
注意:发送事务消息时不能使用普通的 DefaultMQProducer,必须使用专门的 TransactionMQProducer,并绑定我们上面写好的监听器。
import org.apache.rocketmq.client.producer.TransactionMQProducer;
import org.apache.rocketmq.client.producer.TransactionSendResult;
import org.apache.rocketmq.common.message.Message;
import java.util.concurrent.*;
public class TransactionProducer {
public static void main(String[] args) throws Exception {
// 1. 创建【事务消息】专用的生产者,设置生产者组名
TransactionMQProducer producer = new TransactionMQProducer("tx_order_group");
producer.setNamesrvAddr("127.0.0.1:9876"); // 请替换为你真实的 NameServer 地址
// 2. 事务回查需要线程池支持(Broker 回查时会使用该线程池中的线程)
ExecutorService executorService = new ThreadPoolExecutor(
2, 5, 100, TimeUnit.SECONDS,
new ArrayBlockingQueue<Runnable>(2000),
r -> new Thread(r, "client-transaction-msg-check-thread")
);
producer.setExecutorService(executorService);
// 3. 【关键点】把刚才写好的监听器绑定到这个生产者上
TransactionListener listener = new OrderTransactionListener();
producer.setTransactionListener(listener);
// 4. 启动生产者
producer.start();
System.out.println("🚀 事务消息生产者启动成功!");
// ==========================================================
// 5. 模拟发送两条订单消息:订单1002(偶数,会成功)和 订单1003(奇数,会失败)
// ==========================================================
// 测试订单 1002
Message msg1 = new Message("OrderTopic", "Tag_Gift", "给用户发 50 积分".getBytes());
msg1.setKeys("1002"); // 把订单号写进 Key 里,方便监听器去查数据库
System.out.println("\n--- 开始发送订单 1002 的事务消息 ---");
// 【注意】发送方法变成了 sendMessageInTransaction
TransactionSendResult result1 = producer.sendMessageInTransaction(msg1, null);
System.out.println("发送半消息完成,结果状态: " + result1.getSendStatus());
// 测试订单 1003
Message msg2 = new Message("OrderTopic", "Tag_Gift", "给用户发 100 积分".getBytes());
msg2.setKeys("1003");
System.out.println("\n--- 开始发送订单 1003 的事务消息 ---");
TransactionSendResult result2 = producer.sendMessageInTransaction(msg2, null);
System.out.println("发送半消息完成,结果状态: " + result2.getSendStatus());
// 为了防止程序一瞬间关闭看不到【回查】效果,我们让主线程多睡一会儿
Thread.sleep(60000);
producer.shutdown();
}
}
三、 控制台运行结果分析
运行上述代码,你的控制台会打印出如下日志:
🚀 事务消息生产者启动成功!
--- 开始发送订单 1002 的事务消息 ---
【本地事务】收到半消息成功响应,开始在本地数据库下单,订单ID: 1002
【本地事务】正在向数据库写入订单信息...
【本地事务】数据库写入成功!准备向 Broker 发送 Commit
发送半消息完成,结果状态: SEND_OK
--- 开始发送订单 1003 的事务消息 ---
【本地事务】收到半消息成功响应,开始在本地数据库下单,订单ID: 1003
【本地事务】正在向数据库写入订单信息...
【本地事务】库存不足,数据库写入失败!准备向 Broker 发送 Rollback
发送半消息完成,结果状态: SEND_OK
💡 深度总结与观察
- 数据的最终一致性:如果你此时启动消费者去消费
OrderTopic,你会发现消费者只能收到“订单1002”的积分消息。 - 为什么收不到 1003?:因为 1003 对应的本地事务失败了,我们在第四步向 Broker 发送了
Rollback指令,Broker 收到后直接在内存中把这条“半消息”抹除了,绝不往下游投递。
这就是 RocketMQ 事务消息的全部奥秘!通过 “半消息确认” + “本地事务执行” + “终审确认/反向回查” 的神级组合拳,优雅地解决了分布式环境下核心业务跨服务的一致性痛点。
如果你觉得这篇文章对你有启发,欢迎点赞、收藏、关注!
更多推荐




所有评论(0)