RocketMQ 事务消息机制与实现
问题
RocketMQ 的事务消息是如何实现分布式事务的?和 Kafka 的事务有什么区别?在生产环境中使用时有哪些坑?
分析
分布式事务是微服务架构中最棘手的问题之一。当一次业务操作涉及多个服务(如订单服务、库存服务、支付服务),如何保证数据的一致性?传统的 XA 两阶段提交虽然能保证强一致性,但性能差、实现复杂、容易造成资源锁定。而 RocketMQ 的事务消息提供了一种最终一致性的轻量级方案,核心思想是"消息发送和本地事务要么同时成功,要么同时失败"。
为什么需要事务消息?
先看一个典型的电商下单场景:
- 用户在下单页面提交订单
- 订单服务创建订单(写入订单表)
- 通知库存服务扣减库存
- 通知积分服务增加积分
如果直接用普通 MQ 消息,可能出现两种情况:
| 场景 | 本地事务 | 消息发送 | 结果 |
|---|---|---|---|
| 先发消息后执行事务 | 消息已发送 | 执行失败 | 库存扣了但订单没创建,数据不一致 |
| 先执行事务后发消息 | 事务已提交 | 发送失败(网络超时/Broker 宕机) | 订单创建了但库存没扣,超卖 |
| 发消息 + 事务都成功但消费端重复消费 | 事务已提交 | 消息被 Broker 重投 | 库存多扣,少卖了 |
事务消息就是为了解决"发送消息"和"本地事务"之间的原子性问题。注意消费端幂等是"另一层问题",事务消息只保证 Producer 端的原子性,消费端幂等要自己实现。
RocketMQ 事务消息的核心机制
RocketMQ 的事务消息使用 半消息(Half Message) + 事务反查 机制,时序流程如下:
Producer Broker Consumer
│ │ │
│ 1. 发送半消息 │ │
├──────────────────────────►│ │
│ │ 写入 RMQ_SYS_TRANS_HALF_TOPIC
│ │ 消息状态 = PREPARED │
│ 2. 返回半消息结果 │ │
│◄──────────────────────────┤ │
│ │ │
│ 3. 回调 executeLocalTransaction() │
│ └─ 执行本地事务 │ │
│ ├─ 订单表 INSERT ✅ │ │
│ ├─ 库存表扣减 ✅ │ │
│ └─ 积分表插入 ✅ │ │
│ │ │
│ 4. 返回 COMMIT/ROLLBACK │ │
├──────────────────────────►│ │
│ │ COMMIT: 投递到目标 Topic │
│ │ ROLLBACK: 删除半消息 │
│ │ │
│ │ 5. Consumer 消费全消息 │
│ ├──────────────────────────►│
│ │ │
│ ── 如果 Producer 在步骤 3 崩溃 ── │
│ │ │
│ │ 6. 定时扫描(60s 一次) │
│ │ 回调 checkLocalTransaction()
│ │◄──────────────────────────┤
│ │ 7. 返回 COMMIT/ROLLBACK │
│ ├──────────────────────────►│
│ │ │这个机制的关键在于半消息:消息先发送到 MQ,但处于"半可见"状态,Consumer 是看不到的。只有本地事务确认成功后,消息才会变成"全可见"状态,投递到目标 Topic。
半消息的存储实现
半消息写入 RocketMQ 的内部系统 Topic:RMQ_SYS_TRANS_HALF_TOPIC。这个 Topic 有 1 个 Queue(默认配置),所有事务消息的半消息都写入这个 Queue。
一条半消息的 CommitLog 记录结构如下:
| 字段 | 内容 | 说明 |
|---|---|---|
| msgId | 系统生成的全局唯一 ID | 用于反查时定位 |
| origTopic | 用户指定的目标 Topic | 如 "order-tx-topic" |
| queueId | 目标 Topic 的 Queue ID | 提交时投递到正确分区 |
| queueOffset | 目标 Queue 的 Offset | 提交时消息顺序 |
| body | 业务消息体 | 序列化后的业务数据 |
| transactionState | COMMIT/ROLLBACK/PREPARED | 默认 PREPARED |
| preparedTransactionOffset | 半消息在 CommitLog 中的偏移 | 提交时回填 |
事务提交时,RocketMQ 从 RMQ_SYS_TRANS_HALF_TOPIC 删除半消息(标记为已提交),同时将消息投递到 origTopic 的 queueId 对应 Queue。这种设计的好处是:对目标 Topic 的 ConsumeQueue 索引没有侵入性,正常消费逻辑无需感知事务消息的存在。代价是写入时多了一次 IO(半消息 Topic 的 CommitLog 写入),提交时又多了两次 IO(删除半消息记录 + 写入目标 Topic 的 CommitLog)。
事务反查的详细机制
反查是 RocketMQ 事务消息最核心的"保底"机制。当 Producer 在提交半消息后崩溃(或者网络分区导致超时),RocketMQ 无法确定本地事务的状态,就会主动回调 Producer 的 checkLocalTransaction() 方法获取事务状态。
反查的全流程时序:
TransactionMessageCheckService 线程
│
│ 每 60s 扫描一次
├── 扫描 RMQ_SYS_TRANS_HALF_TOPIC 的 ConsumeQueue
│
├── 发现一条 PREPARED 状态的半消息
│ └─ 检查是否已超过 transactionTimeOut(默认 6s)
│
├── 检查反查次数是否超过 transactionCheckMax(默认 15 次)
│ ├─ 未超限 → 发送反查请求到 Producer
│ └─ 已超限 → 将消息转移到 DLQ(死信队列)
│
├── Producer 收到反查请求
│ └─ 调用 checkLocalTransaction()
│ ├─ COMMIT_MESSAGE → Broker 提交消息
│ ├─ ROLLBACK_MESSAGE → Broker 回滚消息
│ └─ UNKNOWN → Broker 等待下次反查(60s 后)
│
└── 如果反查请求发送失败(Producer 仍不可达)
└─ 等待下次扫描周期(60s 后再试)反查的触发条件:
- 半消息在
RMQ_SYS_TRANS_HALF_TOPIC中存活超过transactionTimeOut(默认 6s) - 该半消息的反查次数 ≤
transactionCheckMax(默认 15 次) - Broker 端的
TransactionMessageCheckService线程正常运行
反查超时的兜底处理:如果反查达到 15 次仍无法确定事务状态,消息会被转移到死信队列。此时需要人工介入,通过 rmqadmin 工具手动提交或回滚:
# 手动提交事务消息
mqadmin commitTransaction -n 127.0.0.1:9876 -t RMQ_SYS_TRANS_HALF_TOPIC -i msgId
# 手动回滚事务消息
mqadmin rollbackTransaction -n 127.0.0.1:9876 -t RMQ_SYS_TRANS_HALF_TOPIC -i msgId生产事故案例:反查阻塞导致事务消息堆积
背景:某电商平台促销活动期间,订单服务突然响应变慢,订单创建失败率升高。
排查过程:
- 监控看到 RocketMQ Broker 的
RMQ_SYS_TRANS_HALF_TOPIC消息数从正常几百条飙升到 10 万+ - 查看 Broker 日志,发现大量
checkTransaction failed的 WARN 日志 - 进一步排查发现,事务反查调用
checkLocalTransaction()时,订单服务查询数据库的 SQL 走了慢查询(全表扫描) - 反查响应超时,每次都返回 UNKNOWN,导致 Broker 不断重试,进一步加重了 Broker 的扫描线程负载
根因:checkLocalTransaction() 反查接口没有做索引优化,SELECT 用 order_id 查但表没有索引,每次反查耗时 2-3 秒。Broker 的 TransactionMessageCheckService 单线程扫描,反查超时导致后续半消息排队等待,恶性循环。
修复:
- 给
order_id字段加唯一索引,反查耗时降到 1ms - 在
checkLocalTransaction()中增加缓存,30s 内查过的订单直接返回 - 调整
transactionCheckInterval从 60s 降到 30s,加速反查周期 - 手动处理了已经堆积的 10 万+半消息:通过
mqadmin批量查询数据库状态后提交
代码示例
事务消息发送端(修正版)
@Component
public class OrderTransactionProducer {
@Autowired
private RocketMQTemplate rocketMQTemplate;
@Autowired
private OrderService orderService;
// 事务消息发送
public void createOrder(OrderDTO orderDTO) {
// 构建消息,orderId 放入 header 供反查时使用
Message<OrderDTO> message = MessageBuilder
.withPayload(orderDTO)
.setHeader("orderId", orderDTO.getOrderId())
.build();
TransactionSendResult result = rocketMQTemplate.sendMessageInTransaction(
"order-tx-producer-group", // producer group,需与 listener 保持一致
"order-tx-topic", // 目标 Topic
message, // 消息内容
orderDTO.getOrderId() // 事务参数,传给 executeLocalTransaction 的 arg
);
if (result.getLocalTransactionState() == LocalTransactionState.COMMIT_MESSAGE) {
log.info("订单事务消息提交成功: orderId={}", orderDTO.getOrderId());
} else if (result.getLocalTransactionState() == LocalTransactionState.ROLLBACK_MESSAGE) {
log.warn("订单事务消息回滚: orderId={}", orderDTO.getOrderId());
} else {
log.warn("订单事务状态未知,等待反查: orderId={}", orderDTO.getOrderId());
}
}
// 本地事务执行器 + 反查
@RocketMQTransactionListener(txProducerGroup = "order-tx-producer-group")
public class OrderTransactionListener implements RocketMQLocalTransactionListener {
@Override
@Transactional
public RocketMQLocalTransactionState executeLocalTransaction(Message msg, Object arg) {
String orderId = (String) arg;
try {
// 解析消息体
OrderDTO orderDTO = (OrderDTO) ((RocketMQLocalTransactionMessage) msg).getPayload();
// 执行本地事务:创建订单 + 扣减本地库存 + 增加积分
orderService.createOrder(orderDTO);
log.info("本地事务执行成功: orderId={}", orderId);
return RocketMQLocalTransactionState.COMMIT;
} catch (DataIntegrityViolationException e) {
// 订单号重复(幂等 INSERT),无需回滚
log.warn("订单已存在,视为成功: orderId={}", orderId);
return RocketMQLocalTransactionState.COMMIT;
} catch (Exception e) {
log.error("本地事务执行失败,回滚消息: orderId={}", orderId, e);
return RocketMQLocalTransactionState.ROLLBACK;
}
}
@Override
public RocketMQLocalTransactionState checkLocalTransaction(Message msg) {
// 从 header 中获取 orderId
// 注意:RocketMQ Spring 的 checkLocalTransaction 中 msg 是 generic Message
// 如果使用 spring-messaging 的 Message,需要通过 HeaderAccessor 获取
String orderId = null;
try {
// 方式一:直接转换(RocketMQ 原生 Message)
if (msg instanceof org.apache.rocketmq.common.message.Message) {
org.apache.rocketmq.common.message.Message extMsg =
(org.apache.rocketmq.common.message.Message) msg;
orderId = extMsg.getProperty("orderId");
}
// 方式二:spring-messaging Message
else if (msg instanceof org.springframework.messaging.Message) {
orderId = (String) ((org.springframework.messaging.Message<?>) msg)
.getHeaders().get("orderId");
}
if (orderId == null) {
log.warn("反查消息中无 orderId");
return RocketMQLocalTransactionState.UNKNOWN;
}
// 查询订单状态(带缓存,30s 过期)
Order order = orderService.getOrderByIdWithCache(orderId);
if (order == null) {
// 订单不存在,可能事务还没提交,返回 UNKNOWN 等下次反查
return RocketMQLocalTransactionState.UNKNOWN;
}
switch (order.getStatus()) {
case CREATED:
case PAID:
return RocketMQLocalTransactionState.COMMIT;
case CANCELLED:
case REFUNDED:
return RocketMQLocalTransactionState.ROLLBACK;
default:
return RocketMQLocalTransactionState.UNKNOWN;
}
} catch (Exception e) {
log.error("反查异常: orderId={}", orderId, e);
// 返回 UNKNOWN 让 Broker 下次重试,不要返回 ROLLBACK 导致误删
return RocketMQLocalTransactionState.UNKNOWN;
}
}
}
}消费端实现(幂等消费 + 业务处理)
@Component
public class OrderConsumer {
@Autowired
private RedisTemplate<String, String> redisTemplate;
@Autowired
private InventoryService inventoryService;
@Autowired
private PointsService pointsService;
@RocketMQMessageListener(
topic = "order-tx-topic",
consumerGroup = "order-consumer-group",
// 消费者线程数,根据业务吞吐量调整
consumeThreadNumber = 20,
// 最大重试次数,默认 16
maxReconsumeTimes = 3
)
public class OrderMessageListener implements RocketMQListener<OrderDTO> {
@Override
public void onMessage(OrderDTO message) {
String orderId = message.getOrderId();
String dedupKey = "order:dedup:" + orderId;
// 幂等性检查:使用 Redis SETNX + 过期时间
// 注意:过期时间要大于业务处理时间 + 重试窗口,避免"真"重复
Boolean isFirst = redisTemplate.opsForValue()
.setIfAbsent(dedupKey, "1", Duration.ofHours(24));
if (Boolean.FALSE.equals(isFirst)) {
log.info("消息已处理过,跳过幂等消费: orderId={}", orderId);
return;
}
try {
// 执行下游业务:扣库存 + 加积分
// 库存服务
inventoryService.deductStock(message.getProductId(), message.getQuantity());
// 积分服务
pointsService.addPoints(message.getUserId(), message.getTotalAmount() / 10);
// 处理成功,业务完成
} catch (Exception e) {
// 处理失败,删除幂等标记,让 RocketMQ 重试
redisTemplate.delete(dedupKey);
log.error("处理订单消息失败,将重试: orderId={}", orderId, e);
throw new RuntimeException("处理订单消息失败", e);
}
}
}
}事务消息配置参数详解
# application.yml
rocketmq:
name-server: 192.168.1.100:9876;192.168.1.101:9876
producer:
group: order-tx-producer-group
# 发送消息超时时间,默认 3000ms,事务消息建议调大
send-message-timeout: 5000
# 失败重试次数,默认 2
retry-times-when-send-failed: 2
# 事务消息超时时间,默认 6000ms
transaction-timeout: 10000
# 最大反查次数,默认 15
transaction-check-max: 10
# 反查间隔,默认 60000ms
transaction-check-interval: 30000# broker.conf 关键配置
# 事务消息反查线程池大小,默认 4,堆消息多时需调大
transactionCheckThreadPoolNums=8
# 半消息 Topic 队列数,默认 1
halfTopicQueueNums=1
# 事务消息反查间隔,默认 60000ms
transactionCheckInterval=30000
# 事务消息超时时间,默认 6000ms
transactionTimeOut=10000
# 最大反查次数,默认 15
transactionCheckMax=10总结
RocketMQ 的事务消息用半消息 + 回调反查机制,实现了 Producer 端本地事务和消息发送的原子性,是一种最终一致性的分布式事务方案。它和 Kafka 的事务有本质区别:Kafka 的事务是跨分区原子写入的 EOS(Exactly-Once Semantics),解决的是流处理中的精确一次语义;RocketMQ 的事务消息解决的是跨系统的分布式事务一致性,通过反查机制保证即使 Producer 崩溃也能恢复事务状态。
使用注意事项
事务超时时间:
transactionTimeOut默认为 6 秒,如果本地事务执行超过该时间,RocketMQ 会开始回调反查。对于执行时间较长的本地事务(如跨库写入、文件上传),需要适当调大这个值。调大后同时调整transactionCheckInterval,避免反查线程空转。反查接口必须幂等且快:
checkLocalTransaction()可能被多次调用,必须保证幂等性。实现方式:以订单 ID 为唯一标识,查询数据库判断事务状态。反查接口必须快速返回(< 100ms),否则会阻塞 Broker 的反查线程,导致大量半消息堆积。Consumer 侧必须幂等:事务消息可能被多次投递(如反查超时后消息被重新投递、Consumer 端处理超时后 Broker 重投),消费端必须做好幂等处理。推荐使用 Redis SETNX 或数据库唯一键做去重表。
半消息 Topic 堆积监控:所有事务消息的半消息都存储在
RMQ_SYS_TRANS_HALF_TOPIC,如果大量事务消息长时间处于"半消息"状态(反查超时或 Producer 长时间不可达),会导致这个 Topic 堆积,影响 Broker 性能。需要配置 Prometheus 告警,阈值建议:半消息数 > 5000 触发告警。反查失败的处理:
transactionCheckMax默认 15 次,超过后消息被转移到 DLQ。需要配置告警通知人工介入,或者通过rmqadmin手动提交/回滚。性能影响:事务消息相比普通消息,多了半消息写入(一次 CommitLog 追加)、反查定时扫描(每 60s 一次)、消息恢复(提交时再写一次 CommitLog)等环节,Pub 端吞吐量大约下降 20-30%。对于高吞吐场景(如日志收集、埋点上报),应该用普通消息。对于低吞吐但对一致性要求高的场景(如订单、支付、资金),值得用事务消息。
与 Kafka 事务的对比
| 特性 | RocketMQ 事务消息 | Kafka 事务 |
|---|---|---|
| 解决的问题 | 分布式事务最终一致性 | 流处理 Exactly-Once |
| 核心机制 | 半消息 + 回调反查 | PID + 序列号 + 事务协调器 |
| 跨系统一致性 | 支持(MQ ↔ 数据库) | 不支持(仅 Kafka 内部原子写入) |
| Consumer 隔离 | 半消息对 Consumer 不可见 | isolation.level=read_committed 隔离事务消息 |
| 性能影响 | Pub 端下降约 20-30% | 生产端增加 1 次 RTT(事务协调器交互) |
| 运维复杂度 | 需关注反查超时和半消息堆积 | 需关注事务协调器健康、日志清理 |
| 适用场景 | 跨系统分布式事务(订单、支付) | 流式 ETL、Kafka Streams 状态一致性 |
| 反查依赖 | 依赖 Producer 服务在线 | 无(事务状态持久化在 Kafka 内部 Topic) |
替代方案:本地消息表
如果不想依赖 MQ 的事务特性,可以考虑更轻量的本地消息表方案:
@Transactional
public void placeOrder(OrderDTO order) {
// 1. 插入订单表
orderDao.insert(order);
// 2. 插入消息表(同一事务,相同的数据库连接)
messageDao.insert(new MessageRecord(
order.getOrderId(),
"order-topic",
JSON.toJSONString(order),
MessageStatus.PENDING
));
}
// 3. 定时任务轮询消息表,发送未发送的消息
@Scheduled(fixedDelay = 5000)
public void sendPendingMessages() {
// 批量拉取 PENDING 状态的消息,每次 100 条
List<MessageRecord> pending = messageDao.findByStatus(
MessageStatus.PENDING,
PageRequest.of(0, 100)
);
for (MessageRecord record : pending) {
try {
SendResult result = rocketMQTemplate.syncSend(
record.getTopic(),
record.getContent()
);
messageDao.updateStatus(record.getId(), MessageStatus.SENT);
log.info("本地消息表消息发送成功: id={}", record.getId());
} catch (Exception e) {
log.error("本地消息表消息发送失败,下次重试: id={}", record.getId(), e);
// 不更新状态,下次定时任务继续发送
}
}
}
// 4. 补偿:监控长时间未发送的消息
@Scheduled(cron = "0 0 2 * * ?") // 每天凌晨 2 点
public void compensatePendingMessages() {
// 查找超过 1 小时仍未发送的消息
List<MessageRecord> timeout = messageDao.findByStatusAndCreateTimeBefore(
MessageStatus.PENDING,
LocalDateTime.now().minusHours(1)
);
for (MessageRecord record : timeout) {
// 发送告警,人工介入检查
alertService.sendAlert("消息发送超时", record);
}
}本地消息表 vs 事务消息选型对比:
| 维度 | 本地消息表 | RocketMQ 事务消息 |
|---|---|---|
| 依赖 | 需要数据库 + 定时任务框架 | 需要 RocketMQ 4.x+ |
| 复杂度 | 需自行实现消息表管理、轮询、补偿 | MQ 内置,开箱即用 |
| 一致性保证 | 数据库本地事务保证 | 半消息 + 反查保证 |
| 实时性 | 依赖轮询间隔(秒级延迟) | 半消息提交后立即投递(毫秒级) |
| 运维成本 | 需要监控消息表大小和堆积 | 需要监控半消息 Topic 和反查线程 |
| 消息可靠性 | 消息表持久化,不丢消息 | 半消息持久化,不丢消息 |
选型建议:
- 如果订单量不大(日均 < 10 万)、团队已有定时任务基础设施 → 本地消息表,更可控
- 如果订单量大、延迟敏感(秒级必须到达)、已经有 RocketMQ 运维经验 → 事务消息,更高效
- 两种方案可以共存:核心交易链路用事务消息,内部补偿任务用本地消息表做兜底