核心是本地事务与消息发送的原子性协同:RocketMQ事务消息通过半消息+事务回查机制保障;本地消息表则依赖同库事务插入消息记录+定时任务补偿投递+幂等消费,二者均需状态持久化、重试机制与幂等设计。

Java 中实现本地事务与消息发送的最终一致性,核心是解决“事务成功但消息没发出去”或“消息发出去了但本地事务回滚”这类数据不一致问题。直接在事务内发消息不可靠,必须借助事务与消息的协同机制。
使用事务消息(如 RocketMQ 的半消息)
RocketMQ 提供了事务消息机制,它把消息发送拆成两个阶段:预提交(Half Message)和确认提交/回滚。生产者执行本地事务后,根据结果向 Broker 发送 Commit 或 Rollback 指令,Broker 再决定是否投递消息。
- 本地事务逻辑写在 executeLocalTransaction 方法中,成功返回 COMMIT_MESSAGE,失败或未知状态返回 UNKNOWN
- Broker 会定期回调 checkLocalTransaction 查询事务状态,确保最终一致性
- 注意:该机制依赖 Broker 的事务回查能力,且仅 RocketMQ 原生支持,Kafka 和 RabbitMQ 不具备此特性
基于本地消息表 + 定时任务补偿
在业务库中建一张 message_log 表,和业务操作同在一个事务中插入消息记录(状态为“待发送”),再由独立线程异步读取并投递消息,成功后更新状态为“已发送”。
在 Java 中初始化和管理阿里云 SDK客户端。包括单例模式、线程安全、endpoint 与 region 配置、VPC 终端节点、同步与异步等。
- 插入消息记录和业务更新必须在同一个 @Transactional 方法里,保证原子性
- 投递服务需幂等处理:消息重复投递不能导致业务重复执行(例如用唯一业务 ID 做去重)
- 定时任务扫描超时未发送或发送失败的消息,重新尝试,直到成功或达到最大重试次数
利用 Seata 的 AT 模式 + 消息中间件适配
Seata 本身不直接管理消息,但可以将消息发送封装成一个“分支事务”,通过自定义 ResourceManager 把发消息动作纳入全局事务生命周期。更常见的是结合本地消息表,在 Seata 的 global transaction 提交后,再触发消息发送。
立即学习“Java免费学习笔记(深入)”;
- 避免在 Seata 的 global transaction 内直接调用 MQ 客户端,否则可能因分支事务超时导致消息误发
- 推荐做法:全局事务成功提交后,监听 TransactionSuccessEvent 或用 Spring 的 @TransactionalEventListener 触发消息投递
- 此时消息发送失败仍需单独补偿,所以依然要配套状态查询和重试机制
使用 Kafka 的事务 API(仅限 Kafka 场景)
Kafka 自身提供 KafkaProducer#beginTransaction / commitTransaction / abortTransaction,可让消息写入与本地数据库操作形成“跨系统事务语义”,但前提是数据库操作也通过 Kafka Connect 或 Debezium 同步到 Kafka —— 实际项目中较少直接混用。
- 更实用的做法是:DB 事务成功 → 发送一条 Kafka 消息 → 消费端通过 exactly-once 语义保证下游只处理一次
- Kafka 的事务仅保证“消息写入的原子性”,不保证 DB 与 Kafka 之间的强一致,仍属最终一致性范畴
- 需要开启 enable.idempotence=true 和配置 transactional.id

















