响应式事务的收尾去哪了:空完成丢掉回调,有值挂死链路

2026-09-27 Aryee 5

事务提交之后再发一条消息、清一次缓存——这个需求在同步代码里是 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 挂在成功分支上,异常信号直接绕过它。于是「事务回滚、错误正常上抛」这条最容易被测到的路,表现得完全正确;剩下两条测试没覆盖的路——空完成、有值完成——一条丢回调,一条吊住线程。

我这次是在排查一批网关超时的时候定位到它的:事务早已提交、数据也在,但响应不返回。回看代码,挂死在成功路径、静默在空完成路径,两个缺陷互相遮掩,而我的测试只断言了值和错误。

三条我会照着做的检查

  1. 凡是「不管有没有值都要做」的动作,不许只挂在 flatMap 上。 收尾、清理、计数、事件登记都算。写这段代码时问一句:Mono<Void>、Mono.justOrEmpty(null)、过滤条件全不命中的 Flux,走到这里会怎样。要挂在完成信号上就用 then、switchIfEmpty、doFinally,或者把两种完成形态并到同一个函数里。
  2. 用 Sink 就负责到底。 创建方发终止信号,整段 drain 配 .timeout(...)。挂死这种问题最怕「没人等它」——没有 timeout,它会一直挂到连接池耗尽才被看见;有 timeout,它当天就在日志里炸出来。
  3. 测试按完成形态分档,不是按返回值分档。 有值完成、空完成、错误完成各一条,每条断言 expectComplete()(用 StepVerifier),而不是只断言值。Mono<Void> 那条最容易省,而它恰好是业务里最常见的那一类方法签名。

一点收尾

这两个缺陷的形状其实很常见:在响应式里,「没有值」是一种正常的完成,不是「什么都没发生」。 按「有值」写出来的收尾代码,对 Mono<Void> 一律不成立;而按「同步集合」的直觉用了一个可以无限等的队列,就会把挂死留在成功路径上。

如果你的项目里也有提交后回调、事件发布、缓存失效这类「成功之后再做一件事」的逻辑,值得拿这三条去 grep 一遍 flatMap。

你手上那套系统,卡在哪一步?

留个手机号和方便的时间,我打过来先听你说现状与约束,再给可执行的判断:这套系统是该继续修、该动哪里,还是干脆重做一版设计更省。

约一次沟通不接纯 UI 外包,只做系统层面的问题。