响应式事务的收尾去哪了:空完成丢掉回调,有值挂死链路
事务提交之后再发一条消息、清一次缓存——这个需求在同步代码里是 TransactionSynchronizationManager 的几个方法,在 Reactor 里没有对应物,得自己做。我在自己的响应式数据访问层里做了:切面拦 @Transactional,往 Reactor Context 里塞一个回调收集器,业务侧用 mono.as(ReactiveTransactionUtil.afterCommit(action)) 登记,方法跑完由切面统一执行。
这套实现里有两个缺陷,症状正好相反:
- 返回
Mono<Void>的方法(也就是所有save、update、delete那类不吐值的)请求正常返回,回调一条都没执行,日志里什么都没有。 - 返回
Mono<T>的方法请求永远不返回,网关超时、连接占满,而事务其实早就提交了。 - 出错的时候两样都不发作,异常干干净净地冒出来。
第三条最容易骗人:它让这段代码看起来是好的。
先看这套东西是怎么搭起来的
收集器是一个单播 Sink:
Sinks.Many<Mono<Void>> afterCommitSink = Sinks.many().unicast().onBackpressureBuffer();
Object result = pjp.proceed();
if (result instanceof Mono<?> mono) {
return mono.flatMap(value -> drainAfterCommitSink(afterCommitSink).then(Mono.just(value)))
.onErrorResume(e -> Mono.error(e))
.contextWrite(ctx -> ctx.put(AFTER_COMMIT_SINK_KEY, afterCommitSink));
}
业务侧登记回调,就是从这个 Context 里把 Sink 取出来往里头发:
public static <T> Function<Mono<T>, Mono<T>> afterCommit(Mono<Void> action) {
return mono -> mono.flatMap(value -> Mono.deferContextual(ctx -> {
Sinks.Many<Mono<Void>> sink = ctx.getOrDefault(AFTER_COMMIT_SINK_KEY, null);
if (sink != null) {
sink.tryEmitNext(action); // 有活跃事务:登记,等切面统一执行
}
return Mono.just(value);
}));
}
选 unicast().onBackpressureBuffer() 是有理由的:订阅者到来之前发出的元素会被缓存,等那唯一一位订阅者取走。业务代码先 tryEmitNext、切面后 asFlux(),靠的就是这一点。
切面里 drainAfterCommitSink 负责收尾:把 Sink 里攒的回调逐个执行完。
第一个坑:Mono<Void> 的正常完成是「空完成」,flatMap 根本不进
上面两段代码有一个共同的写法问题:该做的事全挂在 flatMap 上。
flatMap 是 onNext 的操作符。而 Mono<Void> 压根不发元素——它的「成功」就是一个 onComplete 信号。所以对一个 @Transactional public Mono<Void> saveOrder(...) 来说:
- 业务侧的
mono.flatMap(...)不触发,回调没登记进 Sink; - 切面的
mono.flatMap(...)也不触发,drain 没跑。
链条一路 onComplete 到底,请求正常返回,谁都没报错,提交后的动作就是不做。在同步代码里「方法返回后」是一个天然存在的时点;在 Reactor 里你必须自己找到那个信号,而 onComplete 和 onNext 是两个不同的信号。
修法是给空完成补一条分支,并把两种完成形态的收尾逻辑抽成同一个函数,别让它们各写一遍:
private static Mono<Void> settle(Mono<Void> action, boolean requireTransaction) {
return Mono.deferContextual(ctx -> {
Sinks.Many<Mono<Void>> sink = ctx.getOrDefault(AFTER_COMMIT_SINK_KEY, null);
if (sink != null) {
sink.tryEmitNext(action);
return Mono.empty();
}
if (requireTransaction) {
log.warn("No active reactive transaction for afterCommit callback");
return Mono.empty();
}
return action; // 没事务可依赖时按策略立即执行
});
}
public static <T> Function<Mono<T>, Mono<T>> afterCommit(Mono<Void> action, boolean requireTransaction) {
return mono -> mono.flatMap(value -> settle(action, requireTransaction).then(Mono.just(value)))
.switchIfEmpty(Mono.defer(() -> settle(action, requireTransaction).then(Mono.<T>empty())));
}
切面那侧同样补 switchIfEmpty,让空完成也把 Sink 排空。
第二个坑:Sink 不收尾,asFlux() 永远不终止
补完第一条之后,有返回值的那一路开始挂死——而且它其实一直是挂死的,只是被第一个坑盖住了。
drainAfterCommitSink 原来的形状是:
private Mono<Void> drainAfterCommitSink(Sinks.Many<Mono<Void>> sink) {
return sink.asFlux().flatMap(mono -> mono, 16, 1).then();
}
Sink 是活的。业务方法跑完只是「不会再有回调进来」,但没有任何人把这个信号发出去,asFlux() 就继续等——单播队列可以无限期地等下去。于是 then() 不结束,切面包出来的那个 Mono 不结束,整条请求链路挂住。
要结束就得显式收尾:
private Mono<Void> drainAfterCommitSink(Sinks.Many<Mono<Void>> sink) {
sink.tryEmitComplete();
return sink.asFlux().flatMap(mono -> { /* 执行回调 */ mono }, 16, 1).then();
}
用 Sink 当收集器,谁创建谁负责它的终止信号——tryEmitComplete 或者 tryEmitError,二选一必须有。同步代码里不存在这个问题:你手里是一个 List,遍历完就完了;换成 Sink 之后「遍历完」这件事需要你亲口说。
为什么出错时反而一切正常
两个缺陷都只在顺利的时候发作,而失败路径把它们全遮住了:
.onErrorResume(e -> {
log.error("Transaction failed, skipping after-commit callbacks", e);
return Mono.error(e);
})
drain 挂在成功分支上,异常信号直接绕过它。于是「事务回滚、错误正常上抛」这条最容易被测到的路,表现得完全正确;剩下两条测试没覆盖的路——空完成、有值完成——一条丢回调,一条吊住线程。
我这次是在排查一批网关超时的时候定位到它的:事务早已提交、数据也在,但响应不返回。回看代码,挂死在成功路径、静默在空完成路径,两个缺陷互相遮掩,而我的测试只断言了值和错误。
三条我会照着做的检查
- 凡是「不管有没有值都要做」的动作,不许只挂在
flatMap上。 收尾、清理、计数、事件登记都算。写这段代码时问一句:Mono<Void>、Mono.justOrEmpty(null)、过滤条件全不命中的Flux,走到这里会怎样。要挂在完成信号上就用then、switchIfEmpty、doFinally,或者把两种完成形态并到同一个函数里。 - 用 Sink 就负责到底。 创建方发终止信号,整段 drain 配
.timeout(...)。挂死这种问题最怕「没人等它」——没有 timeout,它会一直挂到连接池耗尽才被看见;有 timeout,它当天就在日志里炸出来。 - 测试按完成形态分档,不是按返回值分档。 有值完成、空完成、错误完成各一条,每条断言
expectComplete()(用StepVerifier),而不是只断言值。Mono<Void>那条最容易省,而它恰好是业务里最常见的那一类方法签名。
一点收尾
这两个缺陷的形状其实很常见:在响应式里,「没有值」是一种正常的完成,不是「什么都没发生」。 按「有值」写出来的收尾代码,对 Mono<Void> 一律不成立;而按「同步集合」的直觉用了一个可以无限等的队列,就会把挂死留在成功路径上。
如果你的项目里也有提交后回调、事件发布、缓存失效这类「成功之后再做一件事」的逻辑,值得拿这三条去 grep 一遍 flatMap。
本文作者:Aryee 发布时间:2026-09-27 01:22
本文出处:ARYEE.cn 固定链接:https://www.aryee.cn/archives/reactive-tx-after-commit-drain.html
转载或引用请连这一行一起带走;文中的代码与结论按当时环境成立,组件升级后请以官方文档为准。