rocketmq事务消息通过“半消息+状态回查”实现最终一致性,需使用transactionmqproducer并实现localtransactionlistener的executelocaltransaction(执行本地事务并返回状态)和checklocaltransaction(幂等查库确认真实状态),消费端须幂等处理,配套手动ack、死信队列、对账监控与事务日志。

Java 分布式事务中用 RocketMQ 事务消息保证一致性,核心是把本地数据库操作和消息发送绑成一个逻辑原子单元,靠“半消息 + 状态回查”机制实现最终一致,不是强一致,但可靠、可落地。
用 TransactionMQProducer 发送半消息
不能用普通生产者,必须用 TransactionMQProducer,并绑定一个 LocalTransactionListener 实现类。这个监听器要重写两个方法:
-
executeLocalTransaction:收到半消息 ACK 后立即执行本地事务(比如扣库存、改订单状态),返回
COMMIT_MESSAGE、ROLLBACK_MESSAGE或UNKNOWN -
checkLocalTransaction:当返回
UNKNOWN或 RocketMQ 长时间没收到确认时,Broker 会主动回调这个方法查本地事务真实状态,必须根据 DB 中的业务记录(如订单表 status 字段)如实返回结果
注意:回查方法里不能只查内存或缓存,必须查持久化后的业务数据;且该方法需幂等、轻量,避免 DB 压力过大。
在 Java 中初始化和管理阿里云 SDK客户端。包括单例模式、线程安全、endpoint 与 region 配置、VPC 终端节点、同步与异步等。
本地事务必须真正落地再响应
在 executeLocalTransaction 方法里,数据库操作要和业务逻辑严格对齐:
- 比如电商下单,先 insert 订单 + insert 订单明细,再更新商品库存,全部成功才返回 COMMIT;任一失败就 rollback 并返回 ROLLBACK
- 不要在 try-catch 里吞异常后还返回 COMMIT —— 这会导致 DB 回滚了但消息发出去了
- 事务提交后,DB 数据必须已刷盘或至少处于可持久化状态,否则回查可能查不到真实结果
消费端必须做幂等处理
RocketMQ 只保证“至少一次投递”,网络抖动、重试、Broker 重启都可能引发重复消息。消费者不能假设消息只来一次:
- 推荐用业务唯一键 + 数据库唯一约束:比如订单号作为消息 key,插入订单表时用
INSERT IGNORE或ON DUPLICATE KEY UPDATE - 或用Redis 记录已处理 msgId:setnx + 过期时间,简单高效,适合高并发场景
- 复杂业务建议引入状态机:如订单状态流转为「待支付 → 已支付 → 已发货」,只有当前状态允许时才执行下一步,非法状态直接忽略
配套保障不能少
事务消息不是开箱即用的银弹,还要配好几项兜底措施:
-
手动 ACK:消费逻辑执行完再调用
messageExt.getAckIndex()确认,别 autoACK -
死信队列开启:配置
maxReconsumeTimes,超限消息进死信 Topic,人工介入排查 - 对账监控:定时比对订单库和物流库的数据差异,发现不一致及时告警或补偿
- 事务日志落库:把每条事务消息的 msgId、业务ID、本地事务状态、时间戳记到一张事务日志表,方便回溯和人工修复
Java免费学习笔记:立即使用
解锁 Java 大师之旅:从入门到精通的终极指南










