Java 基础体系 · 第 70/100 篇。示例统一以 Java 25 LTS 为语言和 JVM 基线;框架示例使用与其兼容的现代稳定版本。
Java Flow 与响应式流:Publisher、Subscriber、背压和协议正确性
java.util.concurrent.Flow 是 Java 标准库对响应式流协议(Reactive Streams)的接口定义。它解决的不是“如何异步执行一个任务”,而是如何让一个数据生产者和一个数据消费者在异步、跨线程、不同处理速度的情况下,仍然按照明确的协议传递数据。
这个协议的核心约束是:
Subscriber必须先收到onSubscribe;Subscriber通过Subscription.request(n)声明自己还能接收多少个元素;Publisher发送的onNext数量不能超过已请求数量;onError或onComplete是终止信号,终止后不能再发送其他信号;- 同一个订阅上的信号必须保持顺序,不能因为并发发送而破坏协议。
因此,响应式流的难点不在接口数量,而在于需求量、并发、取消、缓存和终止状态必须同时保持一致。
Flow API 与响应式流的关系
Java 9 引入了 java.util.concurrent.Flow,Java 25 LTS 中仍然使用这组 API。它包含四个核心接口:
public final class Flow {
public interface Publisher<T> {
void subscribe(Subscriber<? super T> subscriber);
}
public interface Subscriber<T> {
void onSubscribe(Subscription subscription);
void onNext(T item);
void onError(Throwable throwable);
void onComplete();
}
public interface Subscription {
void request(long n);
void cancel();
}
public interface Processor<T, R>
extends Subscriber<T>, Publisher<R> {
}
}
它们分别表示:
Publisher<T>:可以发布T类型数据的组件;Subscriber<T>:接收T类型数据的组件;Subscription:一次具体订阅关系的控制句柄;Processor<T, R>:同时是Subscriber<T>和Publisher<R>,可以接收T并发布R。
Flow 的接口和 Reactive Streams 规范采用相同的核心协议模型。Reactive Streams 还定义了 org.reactivestreams.Publisher 等接口,它们不属于 JDK。第三方响应式库通常提供两者之间的适配器,但 Java 标准库本身并不自动提供所有外部库的适配实现。
Flow 不是一个完整的响应式编程框架。它没有内置 map、filter、线程调度器、重试、合并等高层操作符,只定义了数据流参与者之间必须遵守的底层契约。
一次订阅包含哪些角色
调用:
publisher.subscribe(subscriber);
并不只是注册一个回调。它会建立一条独立的订阅通道,典型流程如下:
sequenceDiagram
participant P as Publisher
participant S as Subscriber
participant Q as Subscription
S->>P: subscribe(S)
P->>S: onSubscribe(Q)
S->>Q: request(2)
P->>S: onNext(item-1)
P->>S: onNext(item-2)
S->>Q: request(1)
P->>S: onNext(item-3)
S->>Q: cancel()
P-->>S: 停止继续发送
这里有一个重要区别:
Publisher表示“可以发布数据的能力”;Subscriber表示“如何处理数据”;Subscription表示“这一次连接的状态”,包括需求量和取消状态。
同一个 Publisher 可以被多个 Subscriber 订阅。每次订阅通常都应有独立的 Subscription,不能把所有订阅者共享成一个全局需求计数器,否则一个订阅者的 request 会错误地影响其他订阅者。
Subscriber 的生命周期
一个合法的订阅生命周期必须从 onSubscribe 开始:
onSubscribe
├── onNext
├── onNext
├── ...
└── onComplete
或者:
onSubscribe
├── onNext
├── ...
└── onError
onComplete 和 onError 都是终止信号。一条订阅只能有一个终止结果:
- 可以正常完成;
- 可以异常结束;
- 可以被取消而不再收到终止回调;
- 不能先
onComplete再onError; - 不能先
onError再onNext。
cancel() 与 onComplete() 不是同一个概念:
onComplete()表示Publisher已经正常结束;cancel()表示Subscriber不再希望继续接收;- 取消后,通常不会要求
Publisher再向该订阅发送onComplete; - 取消是异步协议中的停止请求,已经在其他线程进入调用过程的信号可能存在竞态,因此实现必须正确处理并发停止。
Subscriber 应当把 onSubscribe 保存下来,并通过它进行需求控制和取消:
final class LoggingSubscriber<T> implements Flow.Subscriber<T> {
private Flow.Subscription subscription;
@Override
public void onSubscribe(Flow.Subscription subscription) {
this.subscription = subscription;
subscription.request(1);
}
@Override
public void onNext(T item) {
System.out.println("item = " + item);
subscription.request(1);
}
@Override
public void onError(Throwable throwable) {
throwable.printStackTrace();
}
@Override
public void onComplete() {
System.out.println("complete");
}
}
这段代码采用“一次请求一个”的策略。收到一个元素、处理完成后,再请求下一个元素。它简单地体现了背压,但不一定有最高吞吐量,因为每个元素都需要一次需求更新。
背压:需求量如何约束发送量
背压(backpressure)是消费者反向约束生产者发送速度的机制。
设某个 Subscription 的需求量为 D:
request(n)中的n必须是正数;- 调用
request(n)后,需求量增加n; - 每发送一个
onNext,需求量减少 1; Publisher发送的元素数不能超过需求量;request(Long.MAX_VALUE)通常表示近似无限需求。
可以用如下不变量表示:
已发送但尚未被需求覆盖的元素数 ≤ 已请求但尚未消耗的需求量
更形式化地,设:
R:累计请求量;E:累计发送的onNext数量;C:是否已经取消或终止。
在没有溢出处理的简化模型中,合法状态必须满足:
E ≤ R
每次事件的状态变化为:
request(n), n > 0:
R := R + n
onNext(item):
发送前必须满足 E < R
发送后 E := E + 1
cancel():
C := true
后续不得继续正常发送
onComplete/onError:
进入终止状态
例如,初始状态为:
R = 0, E = 0
调用 request(2):
R = 2, E = 0
发送第一个元素:
R = 2, E = 1
发送第二个元素:
R = 2, E = 2
此时 E == R,继续发送第三个元素就是协议错误。Subscriber 如果还想接收数据,必须再次调用:
subscription.request(1);
随后:
R = 3, E = 2
于是又允许发送一个元素。
需求量实际表示的是“允许继续发送的额度”,不是已经处理完成的元素数。一个 Subscriber 可以提前请求 100 个元素,然后在内部队列中逐个处理。
request(n) 不是“立即发送 n 个元素”
request(n) 是一个需求信号,不是同步调用命令。它不保证:
Publisher立即发送;Publisher一定能发送恰好n个;onNext会在调用request的线程执行;- 所有请求都能最终满足。
例如,Publisher 可能只有三个元素,而 Subscriber 请求了 100 个。它可以发送三个元素,然后调用 onComplete。需求量没有被“填满”的义务,完成信号可以提前结束数据流。
同样,Publisher 也可以暂时没有数据。请求量可以保持为正数,等数据将来产生时再发送。
请求量的溢出处理
如果实现直接执行:
requested += n;
就可能发生 long 溢出。例如:
requested = Long.MAX_VALUE - 1
n = 10
结果会变成负数,进而破坏需求量判断。
常见的规范实现会把需求量饱和到 Long.MAX_VALUE:
static long addCap(long current, long n) {
long next = current + n;
return next < 0L ? Long.MAX_VALUE : next;
}
这段代码只适用于 current 和 n 都已经确认是正数的场景。生产级实现还需要在原子更新或锁保护下执行,否则多个线程同时调用 request 仍可能丢失更新。
Long.MAX_VALUE 通常被当作无界需求。例如:
subscription.request(Long.MAX_VALUE);
这并不意味着内存可以无限增长,也不代表 Publisher 必须永远运行。它只是表示 Subscriber 不再要求 Publisher 按有限额度限制发送。数据源本身仍可能耗尽、失败或被取消。
一个可运行的 Flow 示例
下面的程序使用 JDK 自带的 SubmissionPublisher,展示:
- Publisher 如何发布数据;
- Subscriber 如何在
onSubscribe中请求数据; onNext处理完成后如何继续请求;- Publisher 如何正常关闭;
- Subscriber 如何接收
onComplete。
import java.util.concurrent.CountDownLatch;
import java.util.concurrent.Flow;
import java.util.concurrent.SubmissionPublisher;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
public class FlowDemo {
public static void main(String[] args) throws InterruptedException {
ExecutorService executor =
Executors.newVirtualThreadPerTaskExecutor();
CountDownLatch done = new CountDownLatch(1);
try (SubmissionPublisher<Integer> publisher =
new SubmissionPublisher<>(executor, 2)) {
publisher.subscribe(new Flow.Subscriber<>() {
private Flow.Subscription subscription;
@Override
public void onSubscribe(Flow.Subscription subscription) {
this.subscription = subscription;
// 先只允许发送一个元素
subscription.request(1);
}
@Override
public void onNext(Integer item) {
System.out.println(
Thread.currentThread().getName()
+ " received " + item);
// 模拟消费耗时
try {
Thread.sleep(50);
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
subscription.cancel();
return;
}
// 当前元素处理完成后,再允许发送一个
subscription.request(1);
}
@Override
public void onError(Throwable throwable) {
throwable.printStackTrace();
done.countDown();
}
@Override
public void onComplete() {
System.out.println("complete");
done.countDown();
}
});
for (int i = 1; i <= 5; i++) {
publisher.submit(i);
System.out.println("submitted " + i);
}
// 关闭表示不再提交新元素。
// 已经提交的元素仍会按协议发送完,然后触发 onComplete。
publisher.close();
done.await();
} finally {
executor.close();
}
}
}
可以使用 Java 25 编译和运行:
javac FlowDemo.java
java FlowDemo
输出顺序会受到异步调度影响,但数据顺序应保持为 1 到 5,并最终出现:
complete
submitted 和 received 的相对顺序不固定,因为 SubmissionPublisher 使用异步执行器投递信号。submitted 5 可能先于 received 1 打印。
这里有三个容易混淆的地方。
第一,submit 不等于 onNext
publisher.submit(i) 表示把元素提交给 SubmissionPublisher。它可能先进入内部缓冲区,之后才由异步任务调用 Subscriber 的 onNext。
因此:
submit(1)
不表示:
onNext(1) 已经执行完毕
第二,背压可能表现为提交方阻塞
示例把每个订阅者的缓冲容量设置为 2:
new SubmissionPublisher<>(executor, 2)
Subscriber 每次只请求一个元素,而且每个元素处理 50 毫秒。如果生产者提交速度更快,缓冲区可能被填满。SubmissionPublisher.submit 在这种情况下可能阻塞提交线程,等待订阅者继续消费。
这说明背压并不总是表现为“生产者收到一个显式的暂停回调”。一种实现可以通过阻塞、排队、丢弃、失败或降低上游生成速度来体现背压,具体取决于组件设计。
第三,关闭 Publisher 不是取消 Subscriber
publisher.close();
表示 Publisher 不再接受新的提交,并在已提交数据处理完成后正常结束订阅。
如果 Subscriber 不想继续接收,应调用:
subscription.cancel();
取消是订阅者主动终止这条订阅通道,而不是让整个 Publisher 对所有订阅者关闭。
没有请求时会发生什么
下面的 Subscriber 在 onSubscribe 中不调用 request:
final class NeverRequestSubscriber<T>
implements Flow.Subscriber<T> {
@Override
public void onSubscribe(Flow.Subscription subscription) {
System.out.println("subscribed");
// 故意不 request
}
@Override
public void onNext(T item) {
System.out.println("item = " + item);
}
@Override
public void onError(Throwable throwable) {
throwable.printStackTrace();
}
@Override
public void onComplete() {
System.out.println("complete");
}
}
合法的 Publisher 不应调用它的 onNext。因为需求量仍为 0,发送任何元素都会违反背压约束。
对于 SubmissionPublisher,元素可能暂时留在内部缓冲区,提交方最终因为缓冲区满而受到阻塞。Publisher 不会因为 Subscriber 没有请求就自动获得发送许可。
这也是一个常见误解:
“异步发送”不等于“可以无限发送”。
异步只说明生产和消费可以在不同线程进行;背压仍然限制着合法发送量。
request(0) 和负数请求
以下代码是错误的:
subscription.request(0);
以及:
subscription.request(-1);
需求量必须是正数。标准实现通常会把非法请求转化为订阅错误,例如向 Subscriber 发送:
onError(new IllegalArgumentException(...))
随后不再继续正常发送。
Subscriber 不应把 request(0) 当作“暂时暂停”的操作。Flow 没有定义这个语义。要暂时停止接收,通常有两种设计:
- 不再追加请求,但保留当前订阅;
- 调用
cancel()永久终止订阅。
第一种只适用于 Publisher 能承受未消费需求和缓冲的情况;第二种不可逆。
回调的线程与串行性
Flow 允许异步和并发执行,但单个订阅上的信号不能随意并发交错。
错误实现可能这样发送:
executor.submit(() -> subscriber.onNext(1));
executor.submit(() -> subscriber.onNext(2));
如果两个任务同时运行,Subscriber 可能观察到:
onNext(2)先于onNext(1);- 两个
onNext同时进入; - 一个线程执行
onComplete,另一个线程仍执行onNext。
这会破坏响应式流协议。
合法实现通常需要使用以下机制之一:
- 单线程事件循环;
- 每个订阅一个串行队列;
- 锁;
- 原子状态加无锁 drain loop;
- 等价的串行化调度机制。
需要区分两个概念:
- 信号串行性:同一订阅的回调不能并发或乱序;
- 回调线程固定:回调是否始终在同一个线程执行。
协议通常要求前者,不必要求后者。一个 Publisher 可以在不同线程执行连续的 onNext,只要这些调用在逻辑上严格串行且有序。
另一方面,Subscriber 的实现仍然必须考虑线程安全。虽然同一订阅的信号应当串行,但 request 和 cancel 可能被业务线程、超时线程或回调线程同时调用。保存的状态、计数器和取消标志不能依赖未经保护的普通变量。
onNext 中再次 request 的问题
同步 Publisher 可能在 request(n) 调用栈内直接触发 onNext:
Subscriber.onSubscribe
-> subscription.request(1)
-> Publisher.onNext(item)
-> Subscriber.onNext(item)
-> subscription.request(1)
-> Publisher.onNext(next)
如果数据很多,这种实现会产生深度递归,甚至导致栈溢出。
因此,Publisher 不能简单地把 request 实现成:
requested += n;
while (requested > 0 && hasNext()) {
subscriber.onNext(next());
requested--;
}
然后允许 onNext 内部再次进入同一个循环。更可靠的做法是使用“工作进行中”标志或队列,把重入请求转换为后续 drain:
request(n)
增加需求量
如果当前没有 drain:
标记 drain 进行中
循环发送
onNext 中再次 request(n)
只增加需求量
发现 drain 已经进行中
返回,由外层 drain 继续处理
这类 drain loop 同时解决两个问题:
- 避免回调重入造成递归;
- 让需求量更新和发送过程能够被串行化。
一个需求受控的 Publisher 状态模型
对于单个订阅,可以把状态抽象为:
ACTIVE
├── request(n > 0) -> ACTIVE,增加 demand
├── onNext -> ACTIVE,消耗 demand
├── cancel -> CANCELLED
├── onComplete -> COMPLETED
└── onError -> FAILED
CANCELLED、COMPLETED、FAILED 都是终止状态:
终止状态
└── 后续 request/onNext/终止信号不再产生正常数据行为
但是,状态转换必须处理并发竞争。例如:
- 线程 A 正准备发送
onNext; - 线程 B 同时调用
cancel(); - 线程 A 已经通过检查,线程 B 才完成取消。
如果实现只用普通布尔值检查:
if (!cancelled) {
subscriber.onNext(item);
}
那么检查和发送之间仍有竞态。cancel() 不是一个天然的全局中断点。实现需要通过锁、原子状态或串行队列定义明确的竞态结果。
实际协议通常允许“已经开始进行的信号”完成,但取消后不应继续无界地发送后续信号。不能把取消理解为能够撤销已经进入用户代码的 onNext。
错误信号与异常传播
数据流中的失败应通过:
subscriber.onError(exception);
传播,而不是让 Publisher 在线程中静默抛出异常。
例如,一个计算型 Publisher 处理元素时失败:
try {
R result = transform(item);
subscriber.onNext(result);
} catch (Throwable ex) {
subscriber.onError(ex);
}
不过,生产级实现不能只写这个 try-catch。它还必须保证:
onError只发送一次;onError发送后不再onNext;- 已取消的订阅不再继续发送;
- 其他订阅者是否受影响,要由 Publisher 的语义决定。
一个订阅失败,不一定要求整个多订阅者 Publisher 关闭。通常每条订阅拥有独立状态;但如果底层数据源本身失败,则所有订阅者可能都会收到错误。
Subscriber 的回调也可能抛出异常。例如:
@Override
public void onNext(Integer item) {
throw new RuntimeException("consumer bug");
}
Subscriber 不应依赖“通过抛出异常通知 Publisher”这一方式传递业务失败。业务异常应该由 Subscriber 自己决定如何处理,或者使用更高层框架定义的错误通道。
Publisher 如果希望健壮地处理 Subscriber 回调异常,通常需要捕获异常、取消该订阅并停止后续信号。但具体清理行为属于实现责任,不能把 Subscriber 抛出的异常当作正常的 onError 业务信号。
Publisher、Subscriber 和 Processor 的数据流方向
Processor<T, R> 同时实现两个方向:
上游 Publisher<T>
|
v
Processor<T, R>
|
v
下游 Subscriber<R>
Processor 必须分别处理两套需求关系:
- 下游请求了多少个
R; - Processor 向上游请求了多少个
T。
例如,一个一进一出的映射 Processor:
下游 request(5)
-> Processor 最多向上游 request(5)
-> 上游最多发送 5 个 T
-> Processor 产生最多 5 个 R
-> Processor 发送给下游
如果 Processor 一收到上游数据就继续向上游请求,而不考虑下游需求,就会形成失控缓冲:
上游高速生产
-> Processor 无限接收
-> Processor 内部队列增长
-> 内存压力或延迟增加
因此,Processor 的核心不是简单地在 onNext 中转换对象,而是维护上下游需求的对应关系。对于一进一出的同步变换,常见策略是:
下游需求 = d
Processor 向上游请求 d
收到一个 T
转换为一个 R
发送一个 R
需求各减少一个
但对于过滤器、批处理器和展开器,需求关系并不相等:
filter可能收到 10 个T,只产生 2 个R;buffer可能收到 10 个T,产生 1 个List<T>;flatMap可能收到 1 个T,产生多个R;- 错误重试可能让一个上游元素被多次请求。
所以 Processor 必须根据操作符的基数变化重新设计需求传播,不能盲目把下游 request(n) 原样转发给上游。
背压不能消除所有内存风险
背压约束的是 Publisher 向 Subscriber 发送 onNext 的数量,但它不自动保证整个系统没有缓存。
一个组件可能合法地:
- 接收上游数据;
- 把数据放入自己的队列;
- 暂时不向下游发送;
- 等待下游需求增加。
这仍然可能导致队列无限增长。
因此,需要区分:
协议正确性:
没有超过下游 request 的发送额度
资源安全性:
队列、线程、连接、文件句柄等不会无限增长
一个无限缓冲的实现可以在协议上完全正确,却最终因为内存耗尽而失败。反过来,一个有界队列也不自动保证协议正确,因为它仍可能在没有需求时错误发送 onNext。
当上游速度长期大于下游处理速度时,系统必须选择一种策略:
- 让生产者阻塞;
- 降低生产速度;
- 将数据持久化到磁盘或外部队列;
- 丢弃最新数据;
- 丢弃最旧数据;
- 合并或采样数据;
- 让订阅失败。
这些是系统语义选择,不是 Flow 接口单独规定的统一答案。
SubmissionPublisher 的定位和边界
SubmissionPublisher<T> 是 JDK 提供的一个通用 Publisher 实现,适合把生产者提交的数据异步分发给多个 Subscriber。
它提供了:
SubmissionPublisher<T> publisher =
new SubmissionPublisher<>(executor, maxBufferCapacity);
其中:
executor决定异步投递任务使用的执行器;maxBufferCapacity是每个订阅者缓冲的容量上限;submit提交数据;close正常结束;closeExceptionally以异常结束。
例如:
publisher.closeExceptionally(
new IllegalStateException("source failed"));
Subscriber 随后应收到 onError,而不是 onComplete。
SubmissionPublisher 适合简单的异步发布场景,但它不是完整的响应式处理框架,也不自动提供:
- 复杂操作符组合;
- 数据持久化;
- 跨进程传输;
- 可靠消息确认;
- 任意情况下的零缓存传输;
- 业务级重试和死信处理。
此外,多个 Subscriber 之间的处理速度可能不同。快速 Subscriber 和慢速 Subscriber 通常有各自的缓冲与需求状态,因此慢订阅者的背压不必然阻塞所有订阅者;但如果生产者必须等待最慢订阅者,整体提交仍可能受到影响,具体取决于实现和提交策略。
常见协议错误及其表现
在 onSubscribe 之前发送数据
错误顺序:
onNext(1)
onSubscribe(subscription)
Subscriber 还没有获得 Subscription,无法表达需求或取消订阅。这破坏了生命周期起点。
没有需求却发送数据
subscriber.onSubscribe(subscription);
// Subscriber 尚未 request
subscriber.onNext(item); // 错误
这会让背压失效。测试中常表现为 Subscriber 收到了它从未请求过的元素。
超过需求量发送
Subscriber 请求一个元素:
subscription.request(1);
Publisher 却连续调用:
subscriber.onNext(1);
subscriber.onNext(2);
第二个元素超出了许可额度。
重复发送终止信号
subscriber.onComplete();
subscriber.onComplete();
或者:
subscriber.onError(ex);
subscriber.onComplete();
终止信号必须是一次性的。
终止后继续发送
subscriber.onComplete();
subscriber.onNext(item);
这会让下游无法判断数据流是否真的结束。
并发调用同一 Subscriber
两个线程同时执行:
subscriber.onNext(a);
subscriber.onNext(b);
即使元素数量没有超过需求量,也可能违反信号串行性和顺序要求。
把 request 当作线程池大小
subscription.request(100);
表示允许发送 100 个元素,不表示启动 100 个线程,也不表示 Subscriber 会并行处理 100 个元素。并行度是执行模型问题,需求量是流量控制问题,两者不能混为一谈。
如何诊断一个响应式流实现
遇到重复数据、卡住、内存增长或订阅不结束时,可以沿着一条订阅记录以下事件:
订阅建立
onSubscribe
request(n)
onNext(item)
request(n)
cancel / onError / onComplete
至少应记录:
- 订阅 ID;
- 当前线程;
- 每次
request的参数; - 累计请求量;
- 累计发送量;
- 缓冲区长度;
- 是否取消;
- 是否已经终止。
最有价值的不变量是:
累计发送量 ≤ 累计请求量
如果发现:
累计发送量 > 累计请求量
问题通常位于 Publisher、Processor 或中间操作符,而不是消费业务本身。
如果流“没有任何输出”,应先检查:
onSubscribe是否被调用;- Subscriber 是否调用了
request(n); n是否为正数;- Publisher 是否仍有数据;
- 是否被提前
cancel; - 是否被异常终止;
- 异步执行器是否已经关闭;
- Publisher 是否因为缓冲区满而阻塞提交线程。
如果内存持续增长,应检查:
- Subscriber 是否长期不请求;
- Publisher 是否有无界缓存;
- Processor 是否向上游请求过多;
- 下游需求是否被错误地放大;
- 生产者是否无法被暂停;
- 是否存在已经取消但仍保留数据和回调引用的订阅。
规范保证、实现行为和工程选择
需要把三类结论分开。
规范保证
Flow 响应式流协议要求关注:
onSubscribe先于其他信号;request(n)使用正数;onNext不得超过需求量;- 信号有序且终止信号唯一;
- 取消后应停止继续发送。
这些是协议正确性的基础。
常见实现行为
具体实现可能选择:
- 同步调用还是异步调用
onNext; - 使用线程池、虚拟线程还是事件循环;
- 需求量使用锁还是 CAS;
- 缓冲区满时阻塞、失败还是丢弃;
- Subscriber 异常时如何清理;
- 多个订阅者之间是否相互影响。
这些不能仅由 Flow.Publisher 这个接口名称推断,必须查看具体实现的 API 文档和源代码。
工程语义选择
以下问题没有统一答案:
- 慢消费者应该阻塞生产者还是丢数据;
- 是否需要保序;
- 是否允许重复处理;
- 失败后是否重试;
- 是否需要持久化;
- 是否需要至少一次或至多一次处理;
- 是否可以接受无界需求。
Flow 只提供传输协议,不提供完整的消息交付语义。若系统需要确认、重放、持久化和跨进程容错,还需要消息队列、数据库或更高层的响应式库配合。
结语
Publisher 和 Subscriber 描述的是数据流的两端,Subscription 则是控制这条数据流的状态通道。背压的本质是:消费者先声明需求,生产者只能在需求额度内发送。
一个实现是否“异步”并不能说明它是否正确。真正需要验证的是:
是否先订阅并建立 Subscription?
是否正确处理 request(n)?
onNext 是否超出需求?
信号是否串行、有序?
终止信号是否唯一?
cancel 是否能停止后续发送?
缓冲和上游请求是否有界?
只有同时满足这些条件,Publisher、Subscriber、背压和错误终止才构成一个协议正确的响应式流。
系列导航与关联阅读
- 系列入口:Java 完整学习路线:从 Java 25 语言与 JVM 到 Spring、微服务和生产交付
- 上一篇:Java 25 Scoped Value:上下文传播、不可变绑定与虚拟线程
- 下一篇:JVM 字节码工具:javap、ASM、Byte Buddy、Instrumentation 和验证
官方资料
本文依据 Java、Spring 与相关项目官方文档重新梳理;正文、示例与生产清单由 WR BLOG 编写。

评论
0 条讨论