这是《RocketMQ 零基础实战指南》的独立章节版。本章从概念、实操和生产排查三个视角展开,代码块保留了原书可直接运行的版本。 事务消息解决“本地数据库事务”和“消息发送”之间的原子性问题:先写本地事务,再提交消息;如果本地事务成功但消息提交失败,通过回查补齐状态。它不是分布式事务协议,也不是 XA 事务,而是基于最终一致的事件发送机制。
10.1 问题背景
常见错误写法:
begin transaction
insert order
send message
commit
消息发送成功后数据库回滚,会产生“订单不存在但事件已发出”的脏事件。数据库提交成功后进程崩溃,则可能漏发事件。
正确思路:
half message
-> execute local transaction
-> commit or rollback message
10.2 核心流程
Producer Broker Consumer
|--- half message ---> |
|<-- half ok ---------- |
| |
| execute local tx |
| commit / rollback |
|--- commit ----------> |--- real message ---> |
| |
|<------ checkback ---- |
|--- commit/rollback -> |
状态:
| 状态 | 含义 |
|---|---|
| COMMIT_MESSAGE | 半消息转为真实消息 |
| ROLLBACK_MESSAGE | 删除半消息,不投递 |
| UNKNOW | 稍后回查 |
10.3 生产者实现
public class OrderTransactionListener implements TransactionListener {
private final OrderService orderService;
private final EventRecordRepository eventRecordRepository;
@Override
public LocalTransactionState executeLocalTransaction(Message msg, Object arg) {
OrderCreateCommand command = (OrderCreateCommand) arg;
try {
orderService.createWithEvent(command, msg.getTransactionId());
return LocalTransactionState.COMMIT_MESSAGE;
} catch (Exception e) {
log.error("local transaction failed, orderNo={}", command.orderNo(), e);
return LocalTransactionState.ROLLBACK_MESSAGE;
}
}
@Override
public LocalTransactionState checkLocalTransaction(MessageExt msg) {
String transactionId = msg.getTransactionId();
return eventRecordRepository.findByTransactionId(transactionId)
.map(record -> switch (record.getStatus()) {
case COMMITTED -> LocalTransactionState.COMMIT_MESSAGE;
case ROLLED_BACK -> LocalTransactionState.ROLLBACK_MESSAGE;
case PENDING -> LocalTransactionState.UNKNOW;
})
.orElse(LocalTransactionState.UNKNOW);
}
}
发送:
TransactionMQProducer producer = new TransactionMQProducer("order-tx-producer");
producer.setNamesrvAddr("rocketmq-namesrv:9876");
producer.setTransactionListener(new OrderTransactionListener(orderService, repository));
producer.start();
Message message = new Message("OrderTransactionTopic", "OrderCreated",
event(command).getBytes(StandardCharsets.UTF_8));
TransactionSendResult result = producer.sendMessageInTransaction(message, command);
10.4 本地事务表
事务消息仍推荐与本地事件表配合:
CREATE TABLE outbox_event (
id BIGINT PRIMARY KEY AUTO_INCREMENT,
event_id VARCHAR(64) NOT NULL,
transaction_id VARCHAR(64) NOT NULL,
aggregate_id VARCHAR(64) NOT NULL,
event_type VARCHAR(64) NOT NULL,
status VARCHAR(20) NOT NULL,
payload JSON NOT NULL,
created_at DATETIME(6) NOT NULL,
updated_at DATETIME(6) NOT NULL,
UNIQUE KEY uk_event_id (event_id),
UNIQUE KEY uk_transaction_id (transaction_id),
KEY idx_status_time (status, created_at)
);
本地事务内:
INSERT INTO orders(...);
INSERT INTO outbox_event(...);
回查时以 outbox_event 为准,而不是根据内存对象或不确定的日志判断。
10.5 回查设计
回查必须满足:
- 幂等;
- 可持久查询;
- 返回确定状态;
- 不依赖已经丢失的内存上下文;
- 有超时和最大回查限制;
- 对异常返回
UNKNOW; - 状态最终落库。
推荐回查逻辑:
query outbox by transactionId
found COMMITTED -> COMMIT_MESSAGE
found ROLLED_BACK -> ROLLBACK_MESSAGE
found PENDING -> UNKNOW
not found -> check timeout
before timeout -> UNKNOW
after timeout -> ROLLBACK_MESSAGE 或人工治理
10.6 消费端语义
事务消息只约束生产端提交关系,不保证消费端与下游数据库强一致:
Broker -> Consumer -> downstream DB
消费端仍需要:
- 幂等;
- 重试;
- 死信治理;
- 对账;
- 业务状态机;
- 补偿任务。
如果消费失败,生产端事务不会自动回滚。不要把事务消息当成同步分布式事务。
10.7 与本地消息表对比
| 方案 | 优点 | 限制 |
|---|---|---|
| RocketMQ 事务消息 | 有回查机制,减少轮询 | 依赖客户端和 Broker 能力 |
| 本地消息表 | 简单直观,易审计 | 需要定时扫描和重发 |
| Outbox + CDC | 解耦清晰 | 引入 CDC 组件 |
中大型系统可以组合使用:事务消息负责即时提交,本地表负责审计和补偿。
10.8 生产注意事项
- 生产者实例与监听器正确绑定;
- 回查接口响应要快;
- 回查逻辑必须有持久状态;
UNKNOW不能无限返回;- 记录 half message ID、transactionId 和业务单号;
- 监控半消息数量和回查次数;
- 升级客户端前验证事务 API 兼容性;
- 生产环境演练进程崩溃和主备切换。
10.9 常见错误
| 错误 | 后果 |
|---|---|
| 本地事务内存判断 | 重启后无法回查 |
| 捕获异常后返回 COMMIT | 本地失败但消息投递 |
| 消费失败期待生产事务回滚 | 语义理解错误 |
| 回查无超时 | 半消息长期悬挂 |
| 未记录 transactionId | 无法定位事件 |
| 未做幂等 | 补发导致重复消费 |
本章小结
事务消息通过半消息、本地事务和回查机制,降低数据库提交与消息发送不一致的概率。它提供的是生产端最终一致,不等于全局分布式事务。可靠的落地关键是本地事务表、确定性回查、幂等消费和完善的监控补偿。
思考题
- 半消息为什么对消费者不可见?
- 回查为什么必须查询持久化状态?
- 消费失败时生产端事务是否会回滚?
- 事务消息和本地消息表如何组合?
- 半消息数量持续增长说明什么问题?