《Spring Boot 高级》5.1 响应式类型与背压

从 Mono/Flux 的惰性讲到 subscribe 那一刻的 onSubscribe→request(n)→onNext 次序,用 BaseSubscriber 与 doOnRequest 观察背压协议,再对比 onBackpressureBuffer/Drop/Latest/Error 四种策略的适用场景,并说明 Sinks 如何取代已废弃的 Processor。

本节目标:讲清 Mono / Flux 的语义与惰性、subscribe() 那一刻真正发生的调用次序、Reactive Streams 的 request(n) 如何构成背压,以及 onBackpressure* 系列与 Sinks 各自解决什么问题。
适用版本:Spring Boot 4.1.x(Java 21)

5.1 响应式类型与背压

实战卷解决的是「怎么在 Boot 里开 WebFlux、怎么写一个返回 Mono 的接口」;本节解决的是「Mono / Flux 到底是什么、subscribe 那一刻发生了什么、背压凭什么能生效」。多数人对响应式的印象停在「惰性」「异步」两个词上,但真出问题(接口不返回、数据被吞、内存暴涨)时,说不清是哪一环坏了。

延续第 1 章的图书借阅域:本节用 BorrowingService 查借阅、BorrowingEvent 表示借阅事件流。先看数据类型的语义,再看订阅时发生了什么。

5.1.1 Mono 与 Flux:0/1 与 0/N 的语义

Reactor 只有两个核心发布者类型:

类型元素个数典型场景
Mono<T>0 或 1查单本书、保存一条记录、返回 void 的完成信号
Flux<T>0 到 N列表查询、事件流、SSE、流式读取

两者都是 org.reactivestreams.Publisher 的实现。本机 reactor-core-3.8.7.jar 的字节码里写得很清楚:

public abstract class reactor.core.publisher.Mono<T>
        implements reactor.core.CorePublisher<T>
public abstract class reactor.core.publisher.Flux<T>
        implements reactor.core.CorePublisher<T>

而 reactor.core.CoreSubscriber 与 CorePublisher 都继承自 org.reactivestreams 的对应接口,所以 Reactor 与 Reactive Streams 规范是同一套语义,不是两套。

关键点是「声明」而非「数据」:下面两行代码执行完,不会有任何 I/O 发生。

Mono<Book> book = bookRepository.findById("b-1");       // 只是描述「将来会去查 b-1」
Flux<BorrowingEvent> events = eventPublisher.events("m-1"); // 只是描述「将来会有一串事件」

Mono.just(x)、Flux.range(0, 10) 也一样——它们只是把「怎么产出」这件事记下来。没有订阅者,流水线不会启动。

5.1.2 subscribe() 时到底发生了什么

让流水线运转的唯一动作是订阅。Mono 上有两个签名很值得注意(本机 javap 核实):

public abstract void subscribe(reactor.core.CoreSubscriber<? super T>);
public final void subscribe(org.reactivestreams.Subscriber<? super T>);
public final reactor.core.Disposable subscribe();

第一个是 abstract,由每个操作符子类各自实现;第二个是 final,它把外部传入的 Subscriber 适配成 CoreSubscriber(补上 currentContext() 默认方法)后转调第一个。所以「Mono 的子类」才是真正决定订阅行为的地方。

一次正常订阅,信号次序是固定的:

onSubscribe(Subscription)   ← 发布者先把「订阅句柄」交给订阅者
request(n)                  ← 订阅者通过句柄声明「我要 n 个」
onNext(item) * n            ← 按需推送
onComplete()                ← 或 onError(Throwable)

用 BaseSubscriber 能直接观察到这条链(它本身就是 CoreSubscriber + Subscription):

Flux<BorrowingEvent> events = Flux.just(
        new BorrowingEvent("m-1", "b-1", "BORROW"),
        new BorrowingEvent("m-1", "b-2", "RETURN"));

events.subscribe(new BaseSubscriber<>() {
    @Override
    protected void hookOnSubscribe(Subscription subscription) {
        System.out.println("[onSubscribe] 拿到 Subscription");
        request(1);                 // 先只要 1 个
    }

    @Override
    protected void hookOnNext(BorrowingEvent event) {
        System.out.println("[onNext] " + event);
        request(1);                 // 处理完再要 1 个
    }

    @Override
    protected void hookOnComplete() {
        System.out.println("[onComplete]");
    }
});

BaseSubscriber 提供的钩子(本机核实):hookOnSubscribe(Subscription)、hookOnNext(T)、hookOnComplete()、hookOnError(Throwable)、hookOnCancel()、hookFinally(SignalType),以及 request(long)、requestUnbounded()、cancel()、upstream()、dispose()。

把 request(1) 改成 requestUnbounded(),onNext 就会连续被调用两次而不再等消费端——这就是有没有背压的全部差别。

5.1.3 request(n):背压的协议层

Reactive Streams 规范(本机 reactive-streams-1.0.3.jar 核实)只有四个接口,全部围绕背压:

public interface Publisher<T>    { void subscribe(Subscriber<? super T> s); }
public interface Subscriber<T>   { void onSubscribe(Subscription s); void onNext(T t);
                                   void onError(Throwable t); void onComplete(); }
public interface Subscription    { void request(long n); void cancel(); }
public interface Processor<T,R>  extends Subscriber<T>, Publisher<R> { }

背压的载体就是 Subscription.request(long):消费者用它告诉生产者「我现在最多还能收 n 个」。生产者必须按这个额度推送,超发就违反规范。request(Long.MAX_VALUE) 表示无界(放弃背压)。

观察 request 的调用最省事的办法是 doOnRequest 与 doOnSubscribe(本机核实存在):

Flux.range(1, 100)
    .doOnSubscribe(s -> System.out.println("subscribed"))
    .doOnRequest(n -> System.out.println("request(" + n + ")"))
    .doOnNext(i -> System.out.println("next " + i))
    .limitRate(10)              // 分批预取,而不是一次要 100
    .subscribe();

limitRate(int) 会把下游的一次大 request 拆成多次小预取(默认 75% 补给),这样即使下游写的是 requestUnbounded(),上游也不会一次性把全部数据灌进来。

背压要成立有前提:上游必须是一个「听 request 的」发布者。Flux.range、Flux.fromIterable、数据库的响应式驱动都属于这一类。但热源(事件总线、Flux.create 里手动 emit)根本不理会下游要多少,这时就需要显式的背压策略。

5.1.4 背压策略:onBackpressure* 系列

当上游不遵守 request(n) 时,用 onBackpressure* 在中间插一个「缓冲区 + 丢弃策略」。本机核实的四个操作符:

操作符行为适用场景
onBackpressureBuffer()无界缓冲,绝不丢下游只是短暂变慢,且确信数据量可控
onBackpressureBuffer(int maxSize, BufferOverflowStrategy)有界缓冲 + 溢出策略需要明确内存上限
onBackpressureDrop() / onBackpressureDrop(Consumer<T>)超出额度的元素直接丢实时行情、日志采样——丢几条无所谓
onBackpressureLatest()只保留最新的一个,旧的覆盖状态类信号,只关心「当前值」
onBackpressureError()一旦超出就发 onError绝不能丢数据,宁可失败

BufferOverflowStrategy 枚举(本机核实)只有三个值:ERROR、DROP_LATEST、DROP_OLDEST。注意它与 FluxSink.OverflowStrategy 不是同一个枚举,后者是 Flux.create 里 FluxSink 的溢出策略,值更多:

枚举值
reactor.core.publisher.BufferOverflowStrategyERROR、DROP_LATEST、DROP_OLDEST
reactor.core.publisher.FluxSink$OverflowStrategyIGNORE、ERROR、DROP、LATEST、BUFFER

选择原则很直白:先问「这条数据丢了会怎样」。会丢钱就 onBackpressureError(显式失败);只是展示用就 onBackpressureDrop;表示「最新状态」就用 onBackpressureLatest;只有当你能论证上游峰值可控时才用无界 onBackpressureBuffer()——它是内存泄漏的常见入口。

5.1.5 Sinks 与 Processor 的关系

Reactor 早期用一个叫 Processor 的类型表示「既是订阅者又是发布者」的中转站:EmitterProcessor、DirectProcessor、UnicastProcessor,都继承 FluxProcessor 并实现 org.reactivestreams.Processor。要在代码里手动 onNext 往流水线里推数据,当时只能靠它们。

这个 API 已经过时了。 本机 reactor-core-3.8.7.jar 的核实结果:

  • reactor.core.publisher.Processor 接口已不存在(javap 报「找不到类」,jar 里也没有对应 .class);
  • FluxProcessor 与 EmitterProcessor 等仍在 jar 里(javap 能列出),但已标记废弃,新代码不应使用。

替代品是 Sinks(本机核实:Sinks.one() / Sinks.many() / Sinks.empty() / Sinks.unsafe())。借阅事件广播可以这样写:

Sinks.Many<BorrowingEvent> sink = Sinks.many().multicast().onBackpressureBuffer();
Flux<BorrowingEvent> hot = sink.asFlux();

hot.subscribe(e -> System.out.println("[subscriber-1] " + e));
hot.subscribe(e -> System.out.println("[subscriber-2] " + e));

sink.tryEmitNext(new BorrowingEvent("m-1", "b-1", "BORROW"));
sink.tryEmitNext(new BorrowingEvent("m-1", "b-2", "RETURN"));
sink.tryEmitComplete();

关键类型与语义(全部本机核实):

类型 / 方法说明
Sinks.Many<T>.tryEmitNext(T)非阻塞尝试发送,返回 EmitResult(成功或失败原因)
Sinks.Many<T>.emitNext(T, EmitFailureHandler)失败时交给 handler 决定重试与否
Sinks.Many<T>.asFlux()转成 Flux 暴露给订阅者
Sinks.Many<T>.currentSubscriberCount()当前订阅者数
Sinks.EmitFailureHandler.FAIL_FAST失败即抛,最常用
Sinks.EmitFailureHandler.busyLooping(Duration)忙等重试直到超时
Sinks.MulticastSpec.onBackpressureBuffer()多播 + 缓冲(示例所用)
Sinks.MulticastSpec.directAllOrNothing()直接转发,任一订阅者跟不上就整体失败
Sinks.MulticastSpec.directBestEffort()直接转发,跟不上的订阅者被跳过

directAllOrNothing 与 directBestEffort 的差别,本质上是「多播时以最慢的订阅者为准,还是以最快的为准」。前者保证所有订阅者看到完全一致的事件序列(会拖慢整体),后者保证吞吐但允许个别订阅者丢事件。

多线程 emit 时不要用 tryEmitNext——它不做重试,高并发下会返回 FAIL_NON_SERIALIZED。应改用 emitNext(value, EmitFailureHandler.FAIL_FAST) 或 busyLooping。Sinks.unsafe() 则完全放弃串行化保证,只在你能确定「只有一个线程在 emit」时才用。

5.1.6 惰性与副作用:最容易踩的坑

「惰性」只对操作符成立,对方法调用不成立。下面两行的差别是本节最实用的一条:

// 反例:loadFromRemote() 在构造这一行时就同步执行了,跟订阅无关
Mono<Book> bad = Mono.just(loadFromRemote("b-1"));

// 正例:推迟到订阅时、并且由订阅触发才执行
Mono<Book> good = Mono.fromCallable(() -> loadFromRemote("b-1"));

Mono.just(...) 的入参在调用 just 之前就被求值了。同理,Mono.fromFuture(CompletableFuture.supplyAsync(...)) 里的 supplyAsync 也是立即提交的。要真正做到「订阅时才执行」,用 Mono.fromCallable(Callable)、Mono.fromSupplier(Supplier) 或 Mono.defer(Supplier<Mono<T>>)(三者本机均核实存在)。

验证时机最直接的办法是把副作用挪进回调:

Mono<Book> m = Mono.fromCallable(() -> loadFromRemote("b-1"))
        .doOnSubscribe(s -> System.out.println("订阅发生,才轮到远程调用"))
        .doOnNext(b -> System.out.println("拿到 " + b.id()));

System.out.println("构造完成,尚未订阅");
m.subscribe();   // 到这里才会打印「订阅发生…」

5.1.7 本机可以做的验证

以上结论都能在本机复现,命令如下(JDK 21 + reactor-core-3.8.7.jar):

export JAVA_HOME=/tmp/springboot_book/jdk-21.0.12.1+1/Contents/Home
J=/tmp/springboot_book/jars

# 1) 确认 Mono / Flux 的接口与签名
"$JAVA_HOME/bin/javap" -cp "$J/reactor-core-3.8.7.jar" reactor.core.publisher.Mono

# 2) 确认 onBackpressure* 与两个枚举的全部取值
"$JAVA_HOME/bin/javap" -cp "$J/reactor-core-3.8.7.jar" reactor.core.publisher.Flux | grep onBackpressure
"$JAVA_HOME/bin/javap" -cp "$J/reactor-core-3.8.7.jar" reactor.core.publisher.BufferOverflowStrategy

# 3) 确认 Sinks 的四个入口与 EmitFailureHandler
"$JAVA_HOME/bin/javap" -cp "$J/reactor-core-3.8.7.jar" reactor.core.publisher.Sinks
"$JAVA_HOME/bin/javap" -cp "$J/reactor-core-3.8.7.jar" 'reactor.core.publisher.Sinks$Many'

# 4) 确认旧 Processor 接口已不在
"$JAVA_HOME/bin/javap" -cp "$J/reactor-core-3.8.7.jar" reactor.core.publisher.Processor   # 报找不到类

要观察 request(n) 的真实调用序列,把 5.1.2 的 BaseSubscriber 片段写成 main 直接跑即可,不需要 Spring 上下文——Reactor 的行为不依赖 Boot。

5.1.8 知道之后能做什么

定位「接口不返回」。 返回 Mono 的方法若在内部调用了 block(),或把副作用写在了操作符之外,请求就会挂住。排查时先确认「有没有订阅者」:没有订阅,再多的 map / flatMap 都不会执行。

给热源补背压。 事件总线、WebSocket 广播这类热源不理会 request(n),接一个 onBackpressureDrop 或 onBackpressureLatest 就能避免下游变慢时无限堆积。

避免内存暴涨。 见到 onBackpressureBuffer() 无参版本要警惕——它没有上限。改成带 maxSize 与 BufferOverflowStrategy 的版本,把「堆积上限」变成一个显式决策。

小结

  • Mono / Flux 只是「怎么产出」的声明,不订阅就什么都不发生;订阅时才走 onSubscribe → request(n) → onNext* → onComplete。
  • 背压的载体是 Subscription.request(long),前提是上游遵守它;不遵守的热源要用 onBackpressure* 兜底。
  • 四种策略对应四种数据重要性:Buffer(不丢但要限界)、Drop(可丢)、Latest(只要最新)、Error(宁可失败)。
  • Sinks 已取代废弃的 Processor;tryEmitNext 不重试,多线程要用 emitNext + EmitFailureHandler。
  • 惰性只对操作符成立,Mono.just(f()) 里的 f() 在构造时就跑了;要真延迟就用 fromCallable / defer。

下一节把这套语义放回 WebFlux:请求进来后,框架是怎么把 Mono / Flux 一路接成响应的。

阅读导航:上一节:4.3 过滤器、拦截器与异常解析 · 下一节:5.2 WebFlux 请求处理 。

继续阅读

探索更多技术文章

浏览归档,发现更多关于系统设计、工具链和工程实践的内容。

全部文章 返回首页

「java」更多文章

  1. 《Spring Boot 入门》18.3 打包与运行
  2. 《Spring Boot 入门》18.2 实现
  3. 《Spring Boot 入门》18.1 需求与设计