本节目标:讲清
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.BufferOverflowStrategy | ERROR、DROP_LATEST、DROP_OLDEST |
reactor.core.publisher.FluxSink$OverflowStrategy | IGNORE、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 请求处理 。
继续阅读
探索更多技术文章
浏览归档,发现更多关于系统设计、工具链和工程实践的内容。