RocketMQ 系列文章(进阶篇第 2 篇):事务消息深度解析与分布式事务落地实践
前言:分布式系统中的“数据一致性”痛点
在上一篇中,我们已经掌握了 RocketMQ 高可用集群的部署与运维,解决了生产环境中消息队列的稳定性和并发瓶颈问题。但在分布式系统中,仅保证消息的可靠传输还不够——当业务操作与消息发送存在依赖关系时,很容易出现“业务执行成功、消息发送失败”或“消息发送成功、业务执行失败”的不一致问题。
例如:电商场景中,用户下单后,需要完成“扣减库存”和“发送订单创建消息”两个操作。若扣减库存成功,但消息发送失败,会导致下游服务(如物流、支付)无法感知订单信息,出现数据断层;若消息发送成功,但扣减库存失败,会导致下游服务误处理,出现超卖、重复履约等问题。
为解决这类分布式事务问题,RocketMQ 提供了事务消息特性,基于“两阶段提交”思想,实现业务操作与消息发送的原子性,确保分布式系统的数据一致性。本篇将从原理、特性、实操三个维度,深度解析 RocketMQ 事务消息,并结合电商下单场景,手把手教你落地分布式事务解决方案。
前置要求:
已掌握 RocketMQ 基础 API 开发(生产者、消费者实现),了解 Topic、Tag、Group 核心概念
已部署 RocketMQ 高可用集群(Dledger 或多主多从模式均可,推荐 Dledger 模式保障消息可靠性)
了解分布式事务基本概念(如两阶段提交、最终一致性),具备简单的 Java 开发能力
一、分布式事务与 RocketMQ 事务消息核心原理
在讲解 RocketMQ 事务消息前,我们先明确核心概念,避免混淆,同时理解事务消息的设计思路——本质是通过“消息预发送+状态确认”,将分布式事务转化为本地事务与消息状态的协同,实现最终一致性。
1.1 核心概念辨析
1.1.1 本地事务与分布式事务
本地事务:单个数据库(或服务)内的事务操作,遵循 ACID 原则(原子性、一致性、隔离性、持久性),例如单库的库存扣减、订单插入。
分布式事务:跨多个数据库(或服务)的事务操作,多个操作需同时成功或同时失败,例如“扣减库存(库存库)+ 生成订单(订单库)+ 发送消息(MQ)”,核心痛点是如何保证多节点操作的原子性。
1.1.2 事务消息的核心定位
RocketMQ 事务消息并非直接解决所有分布式事务问题,而是专注于“业务操作与消息发送”的原子性,即:确保“业务执行成功”与“消息发送成功”绑定,要么都成功,要么都失败,为分布式事务的最终一致性提供基础。
补充:RocketMQ 事务消息采用“最终一致性”模型,而非强一致性——允许短暂的不一致,但最终会通过补偿机制(回查、重试)达到一致,适配大多数分布式业务场景(如电商、物流、支付)。
1.2 RocketMQ 事务消息原理(两阶段提交)
RocketMQ 事务消息的核心是“两阶段提交”+“事务回查”,整体流程分为 5 个步骤,结合示意图理解更清晰(无需画图,用文字精准描述):
-
第一阶段:发送半事务消息生产者向 RocketMQ 发送“半事务消息”(Half Message),Broker 接收消息后,会将消息标记为“暂存状态”,此时消息不会被消费者消费(即使消费者订阅了该 Topic)。同时,Broker 会记录消息的事务状态(未确认),并返回消息 ID 给生产者。关键:半事务消息的核心是“可回滚”——若后续本地事务执行失败,可通过消息 ID 主动删除该消息;若执行成功,可确认消息状态,使其变为“可消费”。
-
第二阶段:执行本地事务生产者收到 Broker 返回的“半事务消息发送成功”响应后,执行本地事务(如扣减库存、生成订单)。本地事务的执行结果只有两种:成功、失败(不考虑超时,超时会由后续回查机制处理)。
-
第三阶段:提交或回滚消息生产者根据本地事务执行结果,向 Broker 发送“消息确认”请求:
-
本地事务执行成功:发送“提交消息”请求,Broker 收到后,将半事务消息标记为“可消费”状态,此时消费者可正常消费该消息。
-
本地事务执行失败:发送“回滚消息”请求,Broker 收到后,删除半事务消息,该消息不会被消费者消费,实现“业务与消息”的原子回滚。
-
第四阶段:事务回查(补偿机制)若生产者在执行本地事务或发送“确认请求”时出现异常(如生产者宕机、网络中断),Broker 会定期(默认 1 分钟)向生产者发送“事务回查”请求,询问本地事务的执行结果。生产者需实现“事务回查接口”,根据消息 ID 查询本地事务状态,然后向 Broker 返回“提交”或“回滚”指令,确保消息状态与本地事务状态一致。
-
第五阶段:消息消费与最终一致性 Broker 确认消息可消费后,消费者正常消费消息并执行下游业务(如物流通知、积分增加)。若消费者消费失败,RocketMQ 会通过重试机制确保消息被消费成功,最终实现整个分布式链路的数据一致性。
核心设计亮点:
-
半事务消息:避免了“消息发送成功、业务执行失败”的问题,暂存状态的消息不会被消费,确保下游服务不会误处理。
-
事务回查:解决了“生产者异常”导致的消息状态不一致问题,通过定期回查,补偿未完成的事务流程。
1.3 事务消息与普通消息、延时消息的区别
| 消息类型 | 核心特性 | 适用场景 |
|---|---|---|
| 普通消息 | 发送即确认,可直接被消费,无事务保障 | 无需与业务绑定,如日志通知、非核心业务推送 |
| 延时消息 | 发送后延迟指定时间才被消费,无事务保障 | 订单超时取消、定时提醒等场景 |
| 事务消息 | 半事务消息暂存,需确认后消费,有事务回查机制 | 分布式事务场景,如订单创建、支付回调、库存扣减 |
二、RocketMQ 事务消息核心特性与配置
在使用 RocketMQ 事务消息前,需先了解其核心特性、支持的配置参数,以及使用限制,避免踩坑。
2.1 核心特性
-
原子性保障:确保“本地事务执行”与“消息发送”原子绑定,要么都成功,要么都失败。
-
事务回查:Broker 自动回查未确认的事务消息,解决生产者异常导致的状态不一致问题,回查频率可配置。
-
最终一致性:通过回查、消息重试机制,确保分布式链路最终数据一致,不保证强一致性。
-
兼容性:支持与 RocketMQ 集群无缝集成,无需额外部署组件,API 易用性高。
2.2 关键配置参数(Broker 端)
事务消息的核心配置在 Broker 端,需修改 broker.conf 文件,主要参数如下(生产环境按需调整):
# 开启事务消息功能(默认开启,无需修改)
transactionMsgEnable = true
# 事务回查最大次数(默认 15 次,超过次数则回滚消息)
transactionCheckMax = 15
# 事务回查时间间隔(默认 60000ms,即 1 分钟)
transactionCheckInterval = 60000
# 事务消息存储目录(默认与普通消息一致,无需单独配置)
storePathTransactionMsg = ${storePathRootDir}/transaction_msg
# 事务状态存储目录(默认与普通消息一致)
storePathTransactionIndex = ${storePathRootDir}/transaction_index
说明:若使用上一篇部署的 Dledger 集群,无需额外修改配置,默认已开启事务消息功能;若需调整回查频率或最大次数,修改上述参数后重启 Broker 即可。
2.3 使用限制(必看)
-
事务消息仅支持 可靠同步发送(syncSend),不支持异步发送(asyncSend)和单向发送(sendOneway),因为需要等待 Broker 确认半事务消息发送成功后,再执行本地事务。
-
事务消息的 Topic 无需特殊配置,与普通消息 Topic 一致,但建议单独创建 Topic 用于事务消息,便于运维和排查。
-
事务回查次数有限制(默认 15 次),若超过最大回查次数,Broker 会自动回滚该消息,需在生产者端做好日志记录,避免数据丢失。
-
生产者需实现事务回查接口,确保能根据消息 ID 查询本地事务状态,否则会导致消息状态异常。
三、实操:RocketMQ 事务消息代码实现(Java 版)
本次实操基于 RocketMQ 4.8.0 版本,结合 SpringBoot 整合方式(最常用),实现“电商下单”场景的事务消息:用户下单 → 执行本地事务(扣减库存、生成订单)→ 发送事务消息 → 下游消费消息(通知物流)。
整体分为 3 部分:生产者(事务消息发送+本地事务执行+回查)、消费者(消息消费)、环境准备。
3.1 环境准备
3.1.1 依赖引入(SpringBoot 项目)
在 pom.xml 中引入 RocketMQ SpringBoot 依赖(版本与 RocketMQ 集群一致,4.8.0):
<dependency>
<groupId>org.apache.rocketmq</groupId>
<artifactId>rocketmq-spring-boot-starter</artifactId>
<version>2.2.0</version>
</dependency>
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-web</artifactId>
</dependency>
<dependency>
<groupId>mysql</groupId>
<artifactId>mysql-connector-java</artifactId>
<version>8.0.30</version>
</dependency>
<dependency>
<groupId>org.mybatis.spring.boot</groupId>
<artifactId>mybatis-spring-boot-starter</artifactId>
<version>2.2.2</version>
</dependency>
3.1.2 配置文件(application.yml)
配置 RocketMQ 生产者、消费者信息,以及数据库信息(用于本地事务执行):
spring:
application:
name: rocketmq-transaction-demo
datasource:
driver-class-name: com.mysql.cj.jdbc.Driver
url: jdbc:mysql://localhost:3306/rocketmq_demo?useUnicode=true&characterEncoding=utf-8&serverTimezone=GMT%2B8
username: root
password: 123456
rocketmq:
name-server: 192.168.1.101:9876;192.168.1.102:9876;192.168.1.103:9876 # 集群 NameServer 地址
producer:
group: transaction-producer-group # 事务生产者组(必须配置)
send-message-timeout: 3000 # 发送超时时间
retry-times-when-send-failed: 2 # 发送失败重试次数
consumer:
group: transaction-consumer-group # 消费者组
message-model: CLUSTERING # 集群消费模式
consume-thread-max: 10 # 最大消费线程数
3.1.3 数据库准备
创建 2 张表:订单表(order)、库存表(stock),用于本地事务执行(扣减库存、生成订单):
-- 库存表
CREATE TABLE `stock` (
`id` int(11) NOT NULL AUTO_INCREMENT,
`product_id` varchar(50) NOT NULL COMMENT '商品ID',
`stock_num` int(11) NOT NULL COMMENT '库存数量',
PRIMARY KEY (`id`)
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COMMENT='库存表';
-- 订单表
CREATE TABLE `order` (
`id` varchar(50) NOT NULL COMMENT '订单ID',
`product_id` varchar(50) NOT NULL COMMENT '商品ID',
`user_id` varchar(50) NOT NULL COMMENT '用户ID',
`order_status` int(1) NOT NULL COMMENT '订单状态:0-待支付,1-已支付,2-已取消',
`create_time` datetime NOT NULL COMMENT '创建时间',
PRIMARY KEY (`id`)
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COMMENT='订单表';
-- 初始化库存数据
INSERT INTO `stock` (`product_id`, `stock_num`) VALUES ('prod001', 100);
3.2 生产者实现(核心:事务消息发送+本地事务+回查)
RocketMQ 事务生产者需实现 TransactionListener 接口,重写两个核心方法:
-
executeLocalTransaction:执行本地事务,返回事务状态(COMMIT_MESSAGE:提交、ROLLBACK_MESSAGE:回滚、UNKNOWN:未知,触发回查)。
-
checkLocalTransaction:事务回查方法,根据消息 ID 查询本地事务状态,返回对应的事务状态。
3.2.1 事务监听器实现
import org.apache.rocketmq.spring.core.RocketMQTemplate;
import org.apache.rocketmq.spring.support.RocketMQHeaders;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.messaging.Message;
import org.springframework.messaging.support.MessageBuilder;
import org.springframework.stereotype.Component;
import org.apache.rocketmq.client.producer.TransactionListener;
import org.apache.rocketmq.client.producer.TransactionSendResult;
import org.apache.rocketmq.common.message.MessageExt;
import javax.annotation.Resource;
import java.util.UUID;
@Component
public class OrderTransactionListener implements TransactionListener {
// 注入 RocketMQ 模板
@Autowired
private RocketMQTemplate rocketMQTemplate;
// 注入库存、订单服务(实际开发中需分层,此处简化)
@Resource
private StockService stockService;
@Resource
private OrderService orderService;
/**
* 执行本地事务
* @param msg 半事务消息
* @param arg 自定义参数(此处传递订单信息)
* @return 事务状态
*/
@Override
public LocalTransactionState executeLocalTransaction(Message msg, Object arg) {
try {
// 1. 解析消息中的订单信息(arg 为自定义传递的订单对象)
OrderDTO orderDTO = (OrderDTO) arg;
// 2. 执行本地事务:扣减库存 + 生成订单(原子操作,需用本地事务注解)
stockService.deductStock(orderDTO.getProductId(), 1); // 扣减1个库存
orderService.createOrder(orderDTO); // 生成订单
// 3. 本地事务执行成功,返回提交消息
return LocalTransactionState.COMMIT_MESSAGE;
} catch (Exception e) {
// 4. 本地事务执行失败,返回回滚消息
e.printStackTrace();
return LocalTransactionState.ROLLBACK_MESSAGE;
}
}
/**
* 事务回查:Broker 定期调用,查询本地事务状态
* @param msg 半事务消息
* @return 事务状态
*/
@Override
public LocalTransactionState checkLocalTransaction(MessageExt msg) {
try {
// 1. 从消息中获取订单ID(消息标签或消息体中携带)
String orderId = new String(msg.getBody());
// 2. 查询本地订单状态:若订单存在,说明本地事务执行成功;否则失败
OrderDTO order = orderService.getOrderById(orderId);
if (order != null) {
// 订单存在,提交消息
return LocalTransactionState.COMMIT_MESSAGE;
} else {
// 订单不存在,回滚消息
return LocalTransactionState.ROLLBACK_MESSAGE;
}
} catch (Exception e) {
// 异常情况,返回UNKNOWN,继续回查
e.printStackTrace();
return LocalTransactionState.UNKNOWN;
}
}
/**
* 发送事务消息(对外提供的接口)
* @param orderDTO 订单信息
* @return 发送结果
*/
public TransactionSendResult sendTransactionMessage(OrderDTO orderDTO) {
// 构建消息:消息体为订单ID,标签为order:create(用于消费者过滤)
Message<String> message = MessageBuilder
.withPayload(orderDTO.getId())
.setHeader(RocketMQHeaders.TRANSACTION_ID, UUID.randomUUID().toString())
.build();
// 发送事务消息:参数(Topic:Tag, 消息, 自定义参数(订单对象))
return rocketMQTemplate.sendMessageInTransaction(
"order_topic:create",
message,
orderDTO
);
}
}
3.2.2 本地事务服务实现(库存+订单)
注意:本地事务(扣减库存+生成订单)需保证原子性,需添加 @Transactional 注解,确保两个操作同时成功或同时失败。
// 库存服务
@Service
public class StockService {
@Autowired
private JdbcTemplate jdbcTemplate;
/**
* 扣减库存(本地事务)
*/
@Transactional(rollbackFor = Exception.class)
public void deductStock(String productId, int num) {
// 1. 查询库存是否充足
Integer stock = jdbcTemplate.queryForObject(
"select stock_num from stock where product_id = ?",
Integer.class,
productId
);
if (stock == null || stock < num) {
throw new RuntimeException("库存不足");
}
// 2. 扣减库存
jdbcTemplate.update(
"update stock set stock_num = stock_num - ? where product_id = ?",
num,
productId
);
}
}
// 订单服务
@Service
public class OrderService {
@Autowired
private JdbcTemplate jdbcTemplate;
/**
* 创建订单(本地事务)
*/
@Transactional(rollbackFor = Exception.class)
public void createOrder(OrderDTO orderDTO) {
// 生成订单(订单ID、商品ID、用户ID、状态、创建时间)
jdbcTemplate.update(
"insert into `order` (id, product_id, user_id, order_status, create_time) values (?, ?, ?, ?, ?)",
orderDTO.getId(),
orderDTO.getProductId(),
orderDTO.getUserId(),
0, // 待支付状态
new Date()
);
}
/**
* 根据订单ID查询订单(用于事务回查)
*/
public OrderDTO getOrderById(String orderId) {
try {
return jdbcTemplate.queryForObject(
"select id, product_id, user_id from `order` where id = ?",
(rs, rowNum) -> {
OrderDTO order = new OrderDTO();
order.setId(rs.getString("id"));
order.setProductId(rs.getString("product_id"));
order.setUserId(rs.getString("user_id"));
return order;
},
orderId
);
} catch (Exception e) {
return null;
}
}
}
// 订单DTO(数据传输对象)
@Data
public class OrderDTO {
private String id; // 订单ID(UUID生成)
private String productId; // 商品ID
private String userId; // 用户ID
}
3.2.3 接口测试(Controller)
编写接口,模拟用户下单,调用事务消息发送方法:
@RestController
@RequestMapping("/order")
public class OrderController {
@Autowired
private OrderTransactionListener transactionListener;
@PostMapping("/create")
public String createOrder(@RequestParam String userId, @RequestParam String productId) {
// 1. 生成订单ID(UUID)
String orderId = UUID.randomUUID().toString().replace("-", "");
// 2. 构建订单DTO
OrderDTO orderDTO = new OrderDTO();
orderDTO.setId(orderId);
orderDTO.setProductId(productId);
orderDTO.setUserId(userId);
// 3. 发送事务消息,执行本地事务
TransactionSendResult result = transactionListener.sendTransactionMessage(orderDTO);
// 4. 返回结果
if (result.getLocalTransactionState() == LocalTransactionState.COMMIT_MESSAGE) {
return "订单创建成功,消息已提交,订单ID:" + orderId;
} else if (result.getLocalTransactionState() == LocalTransactionState.ROLLBACK_MESSAGE) {
return "订单创建失败,消息已回滚";
} else {
return "订单创建中,消息待确认(将触发回查)";
}
}
}
3.3 消费者实现(消息消费+下游业务)
消费者订阅事务消息,消费成功后执行下游业务(如通知物流、发送积分),需保证消费的幂等性(避免重复消费)。
import org.apache.rocketmq.spring.annotation.RocketMQMessageListener;
import org.apache.rocketmq.spring.core.RocketMQListener;
import org.springframework.stereotype.Component;
@Component
// 订阅 order_topic 主题,tag 为 create,消费者组为 transaction-consumer-group
@RocketMQMessageListener(topic = "order_topic", selectorExpression = "create", consumerGroup = "transaction-consumer-group")
public class OrderConsumer implements RocketMQListener<String> {
/**
* 消费消息(订单ID)
* @param orderId 消息体(订单ID)
*/
@Override
public void onMessage(String orderId) {
try {
// 1. 幂等性校验:查询该订单是否已处理过(避免重复消费)
// 实际开发中可使用 Redis 或数据库记录消费状态
boolean isProcessed = checkOrderProcessed(orderId);
if (isProcessed) {
System.out.println("订单" + orderId + "已处理,跳过重复消费");
return;
}
// 2. 执行下游业务:通知物流、增加用户积分等
System.out.println("收到订单创建消息,订单ID:" + orderId);
System.out.println("执行下游业务:通知物流发货、增加用户积分");
// 3. 标记订单已处理(幂等性保障)
markOrderProcessed(orderId);
} catch (Exception e) {
// 消费失败,RocketMQ 会自动重试(默认 16 次)
e.printStackTrace();
throw new RuntimeException("消息消费失败,触发重试");
}
}
/**
* 幂等性校验:模拟查询订单是否已处理
*/
private boolean checkOrderProcessed(String orderId) {
// 实际开发中可查询 Redis 或数据库,此处简化为模拟
return false;
}
/**
* 标记订单已处理
*/
private void markOrderProcessed(String orderId) {
// 实际开发中可存入 Redis 或数据库,此处简化为模拟
System.out.println("订单" + orderId + "标记为已处理");
}
}
3.4 测试验证(核心场景)
启动 SpringBoot 项目、RocketMQ 集群,通过 Postman 调用 /order/create 接口,测试 3 种核心场景,验证事务消息的正确性:
场景1:本地事务执行成功,消息正常提交
-
请求参数:userId=user001,productId=prod001(库存充足)。
-
预期结果:库存扣减成功(prod001 库存变为 99)、订单表新增一条订单、消费者收到消息并执行下游业务。
场景2:本地事务执行失败,消息回滚
-
请求参数:userId=user001,productId=prod001(库存改为 0,模拟库存不足)。
-
预期结果:库存扣减失败、订单表无新增记录、Broker 收到回滚请求,删除半事务消息,消费者无消息消费。
场景3:生产者异常,触发事务回查
-
操作:在 executeLocalTransaction 方法中抛出 UNKNOWN 异常(或模拟生产者宕机后重启)。
-
预期结果:Broker 定期调用 checkLocalTransaction 方法,查询订单状态;若订单存在,提交消息;若不存在,回滚消息。
四、生产环境落地注意事项(避坑指南)
事务消息在生产环境落地时,需重点关注幂等性、异常处理、性能优化等问题,避免出现数据不一致、消息堆积等问题。
4.1 幂等性保障(核心避坑点)
消费者必须保证幂等性——即使消息被重复消费,也不会导致下游业务出现异常(如重复发货、重复积分)。常用实现方式:
-
基于订单ID/消息ID幂等:将已处理的订单ID(或消息ID)存入 Redis(设置过期时间)或数据库,消费前先查询是否已处理。
-
基于业务状态幂等:消费时先查询业务状态(如订单是否已通知物流),若已处理则直接跳过。
4.2 异常处理规范
-
本地事务异常:必须捕获所有异常,返回 ROLLBACK_MESSAGE,避免返回 UNKNOWN(减少回查压力);同时做好日志记录,便于排查问题。
-
回查异常:回查方法需保证无异常,若出现异常,返回 UNKNOWN,让 Broker 继续回查;不可随意返回 ROLLBACK_MESSAGE,避免误回滚。
-
消费异常:消费失败时,无需手动处理,依赖 RocketMQ 重试机制;若重试多次仍失败,可将消息转入死信队列,后续人工处理。
4.3 性能优化建议
-
控制回查频率:根据业务场景调整 transactionCheckInterval 参数,避免回查过于频繁(如核心业务可调整为 30s,非核心业务调整为 1min)。
-
批量处理:若存在大量事务消息,可采用批量发送、批量消费,提升吞吐量;但需注意本地事务批量执行的原子性。
-
资源隔离:事务消息的生产者组、消费者组与普通消息隔离,单独创建 Topic,便于运维和性能监控。
4.4 监控与排查
-
监控事务消息状态:通过 RocketMQ 控制台,查看事务消息的状态(未确认、已提交、已回滚),及时发现异常消息。
-
日志排查:生产者需记录本地事务执行日志、消息发送日志、回查日志;消费者需记录消费日志,便于排查数据不一致问题。
-
死信队列处理:配置事务消息的死信队列(与普通消息一致),对重试多次仍失败的消息,人工排查原因后手动处理。
五、本篇核心总结及下一篇预告
RocketMQ 事务消息基于“两阶段提交+事务回查”,核心解决“本地事务与消息发送”的原子性问题,实现分布式系统的最终一致性。
-
核心流程:发送半事务消息 → 执行本地事务 → 提交/回滚消息 → 事务回查(补偿) → 消费者消费,确保数据一致。
-
实操关键:生产者实现 TransactionListener 接口(本地事务+回查),消费者保证幂等性,Broker 配置回查参数。
-
生产落地:重点关注幂等性、异常处理、性能优化,做好监控与日志排查,避免数据不一致和消息堆积。
下一篇,我们将讲解 RocketMQ 进阶特性——消息过滤与消息回溯,解决“精准消费”和“历史消息查询”问题,进一步提升消息队列的灵活性和实用性。
AtomGit 是由开放原子开源基金会联合 CSDN 等生态伙伴共同推出的新一代开源与人工智能协作平台。平台坚持“开放、中立、公益”的理念,把代码托管、模型共享、数据集托管、智能体开发体验和算力服务整合在一起,为开发者提供从开发、训练到部署的一站式体验。
更多推荐



所有评论(0)