在一些对数据一致性有强需求的场景,可以用 Apache RocketMQ 事务消息来解决,从而保证上下游数据的一致性。
以电商交易场景为例,用户支付订单这一核心操作的同时会涉及到下游物流发货、积分变更、购物车状态清空等多个子系统的变更。当前业务的处理分支包括:
当主分支订单系统状态更新失败后,物流、积分、购物车系统都不应该接收到消息
使用普通消息是做不到的,因为他会直接将消息发送到topic中
而事务消息参考了两阶段提交的原理,
整个事务消息的详细交互流程如下图所示:

@Test
public void sendTrans() throws MQBrokerException, RemotingException, InterruptedException, MQClientException {
// 创建事务消息发送客户端
TransactionMQProducer transProducer = new TransactionMQProducer("test-trans-producer");
transProducer.setNamesrvAddr(RocketMQConfig.NAME_SERVER_ADDR);
// 指定回查事务消息时的线程池
ExecutorService executorService = new ThreadPoolExecutor(2, 5, 100, TimeUnit.SECONDS, new ArrayBlockingQueue<>(2000), new ThreadFactory() {
@Override
public Thread newThread(Runnable r) {
Thread thread = new Thread(r);
thread.setName("client-transaction-msg-check-thread");
return thread;
}
});
transProducer.setExecutorService(executorService);
// 设置事务监听器
transProducer.setTransactionListener(new TransactionListener() {
// 执行本地事务
@Override
public LocalTransactionState executeLocalTransaction(Message message, Object o) {
System.out.println(Thread.currentThread().getName() + ":执行本地事务");
// 触发回查机制
return LocalTransactionState.UNKNOW;
}
// 回查本地事务,如果执行本地事务返回UNKNOW状态或者生产者应用退出导致本地事务未提交任何状态
@Override
public LocalTransactionState checkLocalTransaction(MessageExt messageExt) {
System.out.println(Thread.currentThread().getName() + ":触发事务回查");
// 提交事务
return LocalTransactionState.COMMIT_MESSAGE;
}
});
transProducer.start();
Message message = new Message(RocketMQConfig.TEST_TOPIC, "hello world".getBytes());
// 发送事务消息
SendResult send = transProducer.sendMessageInTransaction(message,null);
System.out.println(send.getSendStatus());
Thread.sleep(Integer.MAX_VALUE);
}
注:需要注意的是事务消息的生产组名称 ProducerGroupName不能随意设置。事务消息有回查机制,回查时Broker端如果发现原始生产者已经崩溃,则会联系同一生产者组的其他生产者实例回查本地事务执行情况以Commit或Rollback半事务消息。