把「领域事件的发布」和「数据库事务」良好协同起来,保证内外系统一致。
TL;DR
- Spring 的
ApplicationEvent默认是同步、单线程、非事务感知的:publishEvent调用后,所有监听器就在当前线程串行跑完。 - 数据库事务是否成功,只有
PlatformTransactionManager.commit()成功调用后才是确定的。在 commit 之前发送事件 = 拿一个尚未确认的结果去做外部副作用操作(ES 索引、缓存、MQ、HTTP 通知……),一旦事务回滚就会导致不同系统间状态不一致。 - 正确做法:把「发事件」这个动作延迟到事务成功提交之后再执行。Spring 提供了一个可以处理这种问题的钩子:
TransactionSynchronization#afterCommit。 @Transactional本质是一段 AOP 拦截:进方法前拿连接、setAutoCommit(false)、把连接绑到ThreadLocal;出方法时根据结果commit/rollback,并在这两个时间点通知所有注册进TransactionSynchronizationManager的回调。
一、一个通用的 EventPublisher 长什么样
以「订单发布 / 内容发布 / 任何领域动作」这类场景为例,需求是事务成功后再发出领域事件。可以封装一个薄薄的组件:
package com.example.common.event;
import lombok.RequiredArgsConstructor;
import org.springframework.context.ApplicationContext;
import org.springframework.stereotype.Component;
import org.springframework.transaction.support.TransactionSynchronizationAdapter;
import org.springframework.transaction.support.TransactionSynchronizationManager;
import java.util.Objects;
/**
* Publishes domain events only after the surrounding transaction commits.
* Without a transaction, events are published immediately.
*/
@Component
@RequiredArgsConstructor
public class DomainEventPublisher {
private final ApplicationContext applicationContext;
public void publish(Object event) {
Objects.requireNonNull(event, "domain event");
if (TransactionSynchronizationManager.isSynchronizationActive()) {
TransactionSynchronizationManager.registerSynchronization(
new TransactionSynchronizationAdapter() {
@Override
public void afterCommit() {
applicationContext.publishEvent(event);
}
}
);
return;
}
applicationContext.publishEvent(event);
}
}
上面代码的实际含义是:
如果当前线程正处在一个 Spring 事务里,就把「发事件」这个动作挂到事务成功提交之后再执行;否则立即发。
二、为什么必须要「等事务提交后再发事件」?
看看下面这一段代码:
@Transactional
public void publishOrder(Long orderId) {
Order o = repo.load(orderId);
o.publish();
repo.save(o); // 1. 写库(此时还没 commit)
eventPublisher.publish(new OrderPublished(orderId)); // 2. 发事件
}
假设 OrderPublished 的监听器里做了两件事:
- 给 ES 建索引 / 刷新 Redis 缓存;
- 调用下游的搜索、推送、通知服务。
在边界场景下会出现以下几类问题:
问题 1:事件已发,事务却回滚了
「方法体没抛异常」≠「事务能 commit 成功」。
publish(...) 执行完之后,事务仍可能被回滚。这确实是可能发生的,并且原因有很多。本方法后续代码抛异常、外层还包着事务、commit() 本身失败(唯一键冲突 / 死锁 / 连接断开)、业务里 setRollbackOnly() 等场景都会导致这个结果。
这里牵扯到 Spring 事务传播(Propagation)与回滚边界的一大堆细节。简单来说:一次
@Transactional方法调用不等于一次数据库事务。真正的
conn.commit()只发生在最外层开启事务的方法返回时,所以内层 service 只是「参与者」,最外层调用方后续任意一步失败都会把当前方法的变更一起回滚。另外commit()本身也可能失败:JPA 攒到 flush 才发 SQL、InnoDB 死锁被选为牺牲者、锁等待超时、连接断开、主备切换……都会让conn.commit()抛出TransactionSystemException。
一旦事务真的回滚,事件却已经发出去了,结果就是:
- ES 已经建了索引;
- 缓存已经被刷成「已发布」;
- 下游服务已经收到「订单 X 已发布」的推送。
外部系统和数据库出现不一致,而且这类问题的排查成本通常极高,往往难以定位。
问题 2:监听器读不到最新数据(当它跳出了当前事务时)
同步 @EventListener + 同线程 + 同事务的场景下,监听器 100% 能读到发布者刚写入还未 commit 的数据。 这是「读自己写」(read-your-own-writes)语义,发布者和监听器共用同一条 Connection、同一个 MySQL session。在任何隔离级别下这个都是成立的,问题在于 Event Listener 如果不在同一个事务中的情况:
- 监听器加了
@Transactional(propagation = REQUIRES_NEW)→ 新事务、新连接、新 session; - 监听器加了
@Async→ 换线程,ThreadLocal里的事务上下文没了; - 监听器调下游服务 / 走只读从库 → 走别的连接池、甚至等主从复制。
这就轮到 MySQL 的隔离级别问题登场了。默认的 REPEATABLE READ 下,另一个 session 根本查不到尚未 commit 的行。如果数据库中 order 的这一行数据是在上层调用还未 commit 的事务中写入的,那么拥有独立事物的当前方法执行 repo.load(...) 会返回 null。
快速回忆一下 InnoDB 四个隔离级别的关键点:主语都是「别的事务」能不能看到当前事务未 commit 的写入。
- READ UNCOMMITTED 允许脏读(几乎不用);
- READ COMMITTED 只能看到已 commit 的最新版本;
- MySQL 默认的 REPEATABLE READ 依赖 MVCC,事务里第一次快照读时生成 read view,之后哪怕别人 commit 了新数据也看不到;
- SERIALIZABLE 读也加锁。
注意「读自己写」是 SQL 标准的基本承诺,同一 session 里自己 INSERT/UPDATE 的行紧接着 SELECT 一定能查到,与隔离级别无关,因为 InnoDB 的 MVCC 会跳过对自己 undo 版本的可见性判断。只有当切换到另一条 Connection(另一个 session)时,才会撞上 REPEATABLE READ 的快照语义。
实践中「监听器换事务 / 换线程 / 调下游」恰恰是最常见的写法(为了避免监听器阻塞主流程),但触发它的根因是「监听器切换到了另一个 session」,不是事件机制本身。
问题 3:外部调用被卷进长事务
监听器里如果做 HTTP、Kafka、ES 这种慢 IO,同步调用会延长事务持有时间,直接加大数据库锁竞争、连接池耗尽的风险。所以正确的做法是:先把库写完提交,再对外做副作用。 afterCommit 正是为这种场景提供的 hook。
三、TransactionSynchronizationManager 的底层原理
Spring 事务这套东西核心是 PlatformTransactionManager + TransactionSynchronizationManager。TransactionSynchronizationManager 是一个纯 ThreadLocal 容器,每个线程维护这么几块东西:
| ThreadLocal | 装的东西 |
|---|---|
resources | 当前线程绑定的资源,比如 DataSource → ConnectionHolder。MyBatis / JdbcTemplate 拿 Connection 就是从这里拿。 |
synchronizations | 一组 TransactionSynchronization 回调(即前文注册进去的那个匿名类)。 |
currentTransactionName / readOnly / isolationLevel / actualTransactionActive | 事务元数据。 |
所以 isSynchronizationActive() 本质就是判断「当前线程的 synchronizations ThreadLocal 是不是被初始化过」,一般是外层事务开启时 AbstractPlatformTransactionManager.prepareSynchronization() 调 initSynchronization() 初始化的。
TransactionSynchronization 提供了完整的钩子:
beforeCommit(readOnly) // commit 之前,还能抛异常让事务回滚
beforeCompletion() // commit/rollback 之前一定会执行
afterCommit() // 只有 commit 成功才执行;抛异常会传播出去
afterCompletion(status) // commit/rollback 结束后一定会执行
关键点:afterCommit 是在事务管理器执行完 doCommit()(真的 flush 到 DB 之后)才被回调的。 时序大致是:
业务方法结束
└─ TransactionInterceptor.commit
├─ triggerBeforeCommit() → 调所有 sync.beforeCommit()
├─ triggerBeforeCompletion()
├─ doCommit() → 真正 conn.commit() ★DB 事务已经完成
├─ triggerAfterCommit() → 调所有 sync.afterCommit() ★事件在这里发出
└─ triggerAfterCompletion()
这解释了两件事:
- 为什么
afterCommit之后发事件是安全的?DB 已经正常写入了。 - 为什么
DomainEventPublisher里不需要传PlatformTransactionManager?它不管事务,只是把回调塞进当前线程的ThreadLocal,事务管理器在 commit 完之后自己会来遍历这些回调。
这就是所谓的 「同步器机制」(synchronization):事务不知道有谁订阅了它,但订阅者能收到它 commit 的通知 ——典型的观察者模式,只不过总线是一个 ThreadLocal。
四、ApplicationContext.publishEvent 又做了什么
Spring 事件本身是同步、单线程、非事务感知的:
publishEvent(event)
└─ ApplicationEventMulticaster.multicastEvent
└─ for each ApplicationListener:
listener.onApplicationEvent(event) // 同线程串行调用
所以默认情况下:
- 事件在发布线程里同步跑完所有监听器;
- 监听器和发布者共享同一个事务(如果有事务的话);
- 监听器抛异常会一路抛回发布者,会导致事务回滚。
Spring 4.2 之后还给了两个更常用的等价物:
| 写法 | 效果 |
|---|---|
@EventListener | 同线程同事务同步调用(跟 publishEvent 直接调等价) |
@TransactionalEventListener(phase = AFTER_COMMIT) | Spring 内部完成了「挂到 afterCommit」的等价操作——如果没在事务里,默认不会执行(可用 fallbackExecution = true 打开) |
@Async @EventListener | 换线程执行监听器 |
DomainEventPublisher 其实是手写版的 @TransactionalEventListener(AFTER_COMMIT),只是把「延迟到 commit 之后」的逻辑放在了发布端而不是监听端。
为什么发布端更好?
- 对监听者透明——监听者写普通
@EventListener就行,不需要知道上游有没有事务。 - 只有一个地方决定时序,业务代码里不会散落各种
phase = AFTER_COMMIT。 - fallback 行为直观——没有事务时立即发,业务不会出现「事件丢失」的情况(
@TransactionalEventListener默认在无事务时会静默丢弃,是一个常见陷阱)。
五、一个 @Transactional 方法框架到底做了什么
这是 Spring Boot 框架中最容易被当成黑盒的注解之一,但实际上展开看它其实是 AOP 代理 + PlatformTransactionManager 的组合。本质上就是我们在手写 SQL 执行的时候手动执行 TART TRANSACTION; 和 COMMIT;。
假设:
@Service
public class OrderService {
@Transactional
public void publishOrder(Long id) { ... }
}
启动时
- Spring 扫到
@Transactional(类或方法)→AnnotationTransactionAttributeSource解析出TransactionAttribute(传播行为、隔离级别、只读、超时、rollbackFor……)。 BeanFactory通过InfrastructureAdvisorAutoProxyCreator给这个 bean 生成一个 CGLIB / JDK 动态代理,切面是TransactionInterceptor。@Autowired注入的OrderService实际上是代理对象。
调用时(框架做的事,一步一步)
proxy.publishOrder(id)
└─ TransactionInterceptor.invoke(MethodInvocation)
├─ 1. 解析 TransactionAttribute(该方法的事务定义)
├─ 2. 通过 PlatformTransactionManager.getTransaction(def)
│ ├─ doGetTransaction() // 拿当前线程绑定的资源(如 ConnectionHolder)
│ ├─ 判断 isExistingTransaction()
│ │ ├─ 有:按 propagation 处理(REQUIRED 直接加入;REQUIRES_NEW 挂起旧的开新的;NESTED 起 savepoint;NOT_SUPPORTED 挂起…)
│ │ │
│ │ └─ 无:doBegin()
│ │ ├─ DataSource.getConnection()
│ │ ├─ conn.setAutoCommit(false) ★关键:手动接管提交
│ │ ├─ 设 isolation / readOnly
│ │ └─ TransactionSynchronizationManager.bindResource(dataSource, connHolder)
│ └─ prepareSynchronization() → initSynchronization() ★之后 publisher 才走 sync 分支
│
├─ 3. try { 调用 publishOrder 方法体 }
│ 方法里 MyBatis/JdbcTemplate 通过 DataSourceUtils.getConnection(ds)
│ → 命中 TSM 里已经绑定的那条 Connection(保证同一事务用同一连接)
│
├─ 4a. 抛异常:completeTransactionAfterThrowing(...)
│ ├─ 判断 rollbackFor / noRollbackFor
│ ├─ triggerBeforeCompletion → doRollback → triggerAfterCompletion(ROLLED_BACK)
│ └─ 清理 ThreadLocal、释放连接
│
└─ 4b. 正常返回:commitTransactionAfterReturning(...)
├─ triggerBeforeCommit
├─ triggerBeforeCompletion
├─ doCommit → conn.commit() ★这里之前 DB 都没真的 commit
├─ triggerAfterCommit ★afterCommit 回调(含事件发布)在这里执行
├─ triggerAfterCompletion(COMMITTED)
└─ 清理 ThreadLocal、释放连接(归还给连接池)
@Transactional 就是在方法前后插了一段 try/catch,进方法前拿一条连接、setAutoCommit(false),把连接绑到 ThreadLocal;出方法时根据结果 commit 或 rollback,并在这两个时间点通知所有注册进 TransactionSynchronizationManager 的回调。
几个常见困惑
- 为什么同一个类里 A 调 B(B 有
@Transactional)不生效? A 直接走了this.B(),绕过了代理,没进TransactionInterceptor,也就没走上面的流程。 - 为什么异步方法里的事务和外层不是同一个?
ThreadLocal换线程就没了,getTransaction拿不到已存在的资源,会新开一个。所以@Async里的@Transactional是独立的。 REQUIRES_NEW为什么真的能”另起炉灶”?doSuspend()会把当前ThreadLocal里的资源、同步器打包保存起来,doBegin()拿一条新的Connection,做完以后再doResume()还原。所以外层的 synchronizations(包括注册的afterCommit)不会在内层 commit 时触发,只会在最外层真正的 commit 时触发——这正是所需的语义。
六、把这些串起来看整个业务链路
一次订单发布大致长这样:
Controller
└─ @Transactional publishOrder(...) ← ①拿连接、setAutoCommit(false)、绑 TL
├─ 领域对象 publish()
├─ repository.save(...) ← ②写库(未 commit)
└─ domainEventPublisher.publish(new OrderPublished(...))
└─ 发现 TL 有 sync → 注册 afterCommit 回调 ← ③先排队,暂时不发送事件
方法返回
← 拦截器进入 commit 流程
├─ doCommit → conn.commit() ← ④DB 真正落地
└─ triggerAfterCommit
└─ 回调触发 → applicationContext.publishEvent(event) ← ⑤此时才发事件
└─ 各 @EventListener 同步执行(建索引 / 发通知 / …)
← 清理 TL、归还连接
如果 ② 之后抛异常:
- 走 rollback 分支,
doRollback → conn.rollback(); afterCommit不会触发(只有afterCompletion(ROLLED_BACK)会触发);- 事件从头到尾没发出去 → 外部系统零副作用。
这就是这个 20 多行小类想解决的核心问题:用 Spring 事务同步器把「领域事件的对外发布」和「数据库事务的最终成败」绑成同一件事,保证内外一致。
七、其他优化场景
-
异步化:
afterCommit里当前实现是同步发事件,监听器仍在事务线程跑(只是事务已提交、连接已归还)。慢监听器可以叠加@Async或者事件里带上必要字段、监听器自己异步落 MQ。 -
持久化事件(Transactional Outbox):如果对可靠性要求更高(进程 crash 也不能丢事件),
afterCommit之后进程挂了还是会丢。业界方案是把事件作为一张表和业务表在同一个事务里写,再由 CDC / 定时器去投递。上面那个 publisher 是”最简可用”版本,能覆盖 95% 的场景。 -
等价写法:如果不希望手写 publisher,可以直接使用
@TransactionalEventListener(phase = AFTER_COMMIT, fallbackExecution = true)语义几乎一样。团队选择手写 publisher 通常是为了在无事务时也能发(fallback 明确)+ 对监听者零侵入。
理解了 TransactionSynchronizationManager 是 ThreadLocal + Connection + 回调列表 这个模型,与 Spring 事务相关的问题从此就不再是黑盒了。