分布式事务消息方案
分布式事务是微服务最难的一环。RocketMQ 事务消息用"半消息 + 回查"在消息层解决最终一致。本文从源码拆解实现机制,并给出 Spring Boot 整合方案。
事务消息解决的问题
经典问题
下单服务要完成两个操作:
1. 写订单表(本地数据库)
2. 发消息"订单已创建"(MQ)
问题:
├─ 先写库后发消息:发消息失败 → 数据库有数据但没通知
└─ 先发消息后写库:写库失败 → 消息已发但业务没成传统方案 vs 事务消息
本地消息表:业务 + 消息同库同事务,靠定时轮询补发
事务消息: RocketMQ 半消息 + 回查,Broker 保证最终一致事务消息整体流程
事务消息流程(三个阶段):
1. 发送半消息(Half Message)→ Broker 暂存(对消费者不可见)
2. 执行本地事务 → 返回结果
├─ COMMIT → Broker 确认消息(对消费者可见)
├─ ROLLBACK → Broker 丢弃消息
└─ UNKNOW → Broker 回查
3. 回查(Check):Broker 定期问生产者"本地事务到底成没成"生产者 Broker 消费者
│ │ │
├─ 发送半消息 ─────────────▶ │(暂存,不可见) │
│ │ │
├─ 执行本地事务 │ │
├─ 提交/回滚/未知 ─────────▶ │ │
│ (commit/rollback) │ │
│ ├─ commit → 消息可见 ────────▶ │ 消费
│ └─ rollback → 丢弃 │
│ │ │
│◀───── 回查(未确认时)──── │ │
└─ 返回本地事务状态 ────────▶ │ │半消息机制(Half Message)
什么是半消息
半消息:已发送到 Broker、但对消费者不可见的消息
用途:等待生产者确认"本地事务是否成功"Broker 端处理
java
// Broker 接收半消息 → 放入半消息队列(特殊 Topic)
// HalfMessageQueue:RMQ_SYS_TRANS_HALF_TOPIC
发送半消息:
1. 消息写入 Half Topic(RMQ_SYS_TRANS_HALF_TOPIC)
2. 消息属性标记:真实 Topic 存入消息属性
3. 消费者查询不到(不在真实 Topic)
提交(Commit):
1. 从 Half Topic 取出消息
2. 按真实 Topic 重新写入
3. 消费者可见
回滚(Rollback):
1. 从 Half Topic 取出消息
2. 丢弃(写入系统删除 Topic 或不写)半消息的存储
Half Topic:RMQ_SYS_TRANS_HALF_TOPIC(系统内部)
真实消息被保存为:
Topic = RMQ_SYS_TRANS_HALF_TOPIC
消息属性中记录 原Topic + 原Queue
Commit 时:
从 Half 队列取出 → 恢复到原 Topic/Queue本地事务执行器
生产者侧两个回调
java
// 1. 本地事务执行器
public interface LocalTransactionExecuter {
LocalTransactionState executeLocalTransactionBranch(
Message msg, Object arg); // 执行本地事务,返回状态
}
// 2. 回查回调
public interface TransactionListener {
// 提交半消息时执行本地事务
LocalTransactionState executeLocalTransaction(Message msg, Object arg);
// 回查:Broker 不确定时询问
LocalTransactionState checkLocalTransaction(MessageExt msg);
}事务状态
java
// org.apache.rocketmq.client.producer.LocalTransactionState
public enum LocalTransactionState {
COMMIT_MESSAGE, // 提交:消息可见
ROLLBACK_MESSAGE, // 回滚:丢弃消息
UNKNOW; // 未知:等 Broker 回查
}源码阅读:生产者端
TransactionMQProducer
java
// org.apache.rocketmq.client.producer.TransactionMQProducer
public class TransactionMQProducer extends DefaultMQProducer {
private TransactionListener transactionListener; // 事务监听器
// 发送事务消息
public TransactionSendResult sendMessageInTransaction(
Message msg, Object arg) throws MQClientException {
// 1. 发送半消息(消息会带事务标识)
SendResult sendResult = this.defaultMQProducerImpl.sendMessage(msg);
// 2. 执行本地事务
LocalTransactionState localTransactionState =
transactionListener.executeLocalTransaction(msg, arg);
// 3. 根据结果提交/回滚
switch (localTransactionState) {
case COMMIT_MESSAGE:
// 提交:让 Broker 放行消息
endTransaction(msg, sendResult, COMMIT);
break;
case ROLLBACK_MESSAGE:
endTransaction(msg, sendResult, ROLLBACK);
break;
case UNKNOW:
// 不做处理:等 Broker 回查
break;
}
return new TransactionSendResult(sendResult, localTransactionState);
}
}回查机制(生产者侧)
java
// Producer 后台线程:定期从 Broker 拉取需要回查的消息
public void checkTransactionState(String addr, MessageExt messageExt) {
// 回调业务方
LocalTransactionState state = transactionListener.checkLocalTransaction(messageExt);
// 根据状态响应 Broker
if (state == COMMIT_MESSAGE) {
endTransaction(messageExt, COMMIT);
} else if (state == ROLLBACK_MESSAGE) {
endTransaction(messageExt, ROLLBACK);
}
// UNKNOW → 下次再查
}源码阅读:Broker 端
回查流程
java
// Broker 端:TransactionMessageCheckService(定时线程)
public class TransactionalMessageCheckService extends ServiceThread {
@Override
public void run() {
while (!this.isStopped()) {
// 定时扫描(默认 1 分钟)
this.onWaitEnd();
// 1. 从 Half Topic 取超时未确认的消息
// 2. 发送回查请求给对应生产者
// 3. 等待生产者响应
// 4. 生产者无响应 → 再次回查(有次数限制)
// 5. 超过次数 → 消息丢弃
}
}
}回查规则
回查时机:半消息超时未确认(默认 6 秒后可查)
回查间隔:默认 60 秒一次
回查次数:默认 15 次(约 15 分钟)
超过次数:消息被丢弃与 Spring Boot 整合
引入依赖
xml
<dependency>
<groupId>org.apache.rocketmq</groupId>
<artifactId>rocketmq-spring-boot-starter</artifactId>
<version>2.2.3</version>
</dependency>配置
yaml
rocketmq:
name-server: 127.0.0.1:9876
producer:
group: order-producer-group生产者:事务消息
java
// 1. 定义事务监听器(本地事务 + 回查)
@Component
public class OrderTransactionListener implements RocketMQLocalTransactionListener {
@Autowired
private OrderService orderService;
// 执行本地事务
@Override
public RocketMQLocalTransactionState executeLocalTransaction(Message msg, Object arg) {
try {
// 解析消息 → 业务参数
OrderMessage orderMsg = parse(msg);
// 执行本地业务(写订单表)
orderService.createOrder(orderMsg);
// 成功 → 提交
return RocketMQLocalTransactionState.COMMIT;
} catch (Exception e) {
// 失败 → 回滚
return RocketMQLocalTransactionState.ROLLBACK;
}
}
// 回查:Broker 不确定本地事务结果时调用
@Override
public RocketMQLocalTransactionState checkLocalTransaction(Message msg) {
OrderMessage orderMsg = parse(msg);
// 查询本地事务是否成功(订单是否存在)
boolean exists = orderService.exists(orderMsg.getOrderId());
return exists
? RocketMQLocalTransactionState.COMMIT
: RocketMQLocalTransactionState.ROLLBACK;
}
}java
// 2. 发送事务消息
@Service
public class OrderService {
@Autowired
private RocketMQTemplate rocketMQTemplate;
@Autowired
private OrderTransactionListener transactionListener;
public void createOrderWithMessage(OrderRequest req) {
// 发事务消息(半消息 + 本地事务 + 回查)
Message msg = MessageBuilder.withPayload(new OrderMessage(req)).build();
rocketMQTemplate.sendMessageInTransaction(
"order-topic",
msg,
req, // 透传给监听器的 arg
transactionListener);
}
}消费者:幂等消费
java
@Component
@RocketMQMessageListener(topic = "order-topic", consumerGroup = "order-consumer-group")
public class OrderConsumer implements RocketMQListener<OrderMessage> {
@Autowired
private InventoryService inventoryService;
@Override
public void onMessage(OrderMessage msg) {
// 幂等:防止消息重复(事务消息保证不丢,但可能重复)
if (duplicateService.isProcessed(msg.getOrderId())) {
log.info("重复消息,跳过: {}", msg.getOrderId());
return;
}
// 扣减库存
inventoryService.deduct(msg.getOrderId(), msg.getProductId(), msg.getQty());
duplicateService.markProcessed(msg.getOrderId());
}
}完整场景:下单 + 扣库存
1. 下单服务:
发送半消息"订单已创建" → Broker 暂存
执行本地事务:写订单表(状态:待扣库存)
→ 提交半消息(消息可见)
2. 库存服务:
订阅 order-topic
消费 → 扣减库存 → 幂等(按订单号)
3. 若步骤 1 崩溃:
Broker 回查下单服务 → 查订单表
→ 存在 → COMMIT;不存在 → ROLLBACK与其他分布式事务方案对比
| 方案 | 一致性 | 性能 | 复杂度 | 适用 |
|---|---|---|---|---|
| XA 两阶段 | 强一致 | 低 | 高 | 短事务、同库 |
| TCC | 最终一致 | 中 | 高(需写 Try/Confirm/Cancel) | 资金类 |
| Saga | 最终一致 | 中 | 中(补偿编排) | 长流程 |
| 事务消息 | 最终一致 | 高 | 低 | 下单通知、异步解耦 |
| 本地消息表 | 最终一致 | 高 | 低 | 通用 |
选型建议:
简单异步解耦 → 事务消息
资金敏感强一致 → TCC/XA
长流程编排 → Saga常见问题
- 事务消息一定不丢吗? 半消息阶段不丢(Broker 持久化 + 回查),但可能重复(消费侧需幂等)。
- 回查超时怎么办? 默认 15 次约 15 分钟,可配置回查次数;业务回查逻辑必须幂等可重入。
- RocketMQ 与 Kafka 事务消息区别? Kafka 事务是分区级原子性,RocketMQ 是业务级半消息+回查,语义更贴近业务。
- 半消息一直 UNKNOW? 检查回查回调是否实现、Broker 回查配置、生产者是否在线。