本节摘要:2015 年 Reactive Streams 1.0 用四个 Java 接口把背压从各框架私有实现升格为协议级义务,并纳入 Java 9 的
java.util.concurrent.Flow。
onError 后为何禁止再 onNextProcessor 在 ETL 链中的桥接角色2013 年 Netflix 推荐引擎在高峰触发 GC 停顿,根因之一是 RxJava 的 onBackpressureBuffer 与 Akka HTTP 对「缓冲区满」解读不一致。Lightbend、Pivotal、Red Hat 等厂商在 GitHub 发起 Reactive Streams Initiative,目标仅是:让 Publisher 与 Subscriber 无歧义对话,不提供实现。
这场危机的本质是隐式契约的分裂。各框架都说自己支持背压,但对「缓冲区满」的语义各执一词:有的视为可恢复的背压信号,有的视为终止性错误,有的静默丢弃。当这些框架在 Netflix 这样的大规模系统里混合部署时,语义分歧变成系统性故障——同一个失败,A 框架认为可以重试,B 框架已经终止链路。Reactive Streams 的回应很朴素:把「发布与订阅之间如何对话」写成一份可读、可实现的规范,每个框架只要遵守它,跨框架协作就不再依赖默契。规范的价值不在发明新机制,而在消除歧义。
| 接口 | 角色 | 关键方法 |
|---|---|---|
Publisher<T> |
数据源 | subscribe(Subscriber) |
Subscriber<T> |
消费者 | onSubscribe/onNext/onError/onComplete |
Subscription |
双向中介 | request(long n), cancel() |
Processor<T,R> |
变换节点 | 同时实现 Pub + Sub |
规范要点(SOURCE 3.1):
request(n) 累计不得超过下游处理能力cancel() 后 Publisher 必须停止发射并释放资源// Java 9 Flow(与 RS 规范对齐) Flow.Publisher<String> pub = ...; pub.subscribe(new Flow.Subscriber<>() { private Flow.Subscription sub; public void onSubscribe(Flow.Subscription s) { sub = s; s.request(1); } public void onNext(String item) { process(item); sub.request(1); } public void onError(Throwable t) { log.error(t); } public void onComplete() { log.info("done"); } });
四接口的方法加起来正好 13 个,这是规范「最小完备」设计的体现:不提供任何「方便但可替代」的 API,每个方法都有不可省略的职责。request 与 cancel 是背压与取消的唯二入口;Processor 允许一个组件同时扮演生产者与消费者,让 ETL 管线可以像搭积木一样把处理段串起来。
| 当前状态 | 合法事件 | 下一状态 |
|---|---|---|
| 未订阅 | subscribe() |
已订阅(待 onSubscribe) |
| 已订阅 | onSubscribe |
可请求 |
| 可请求 | request(n) / onNext |
可请求 |
| 可请求 | cancel() / onError / onComplete |
已终止 |
| 已终止 | 任何信号 | 非法(须忽略或抛错) |
这张状态机表是规范的最简模型。关键在于「已终止」是吸收态:一旦 onError 或 onComplete 到达,或 cancel() 被调用,此后任何 onNext 都是违规。这解释了为什么错误必须作为信号显式传播——如果错误被吞掉,下游会永远停在「可请求」状态空等。这也是为什么第六章会把「错误黑洞」(裸 subscribe() 不传 error consumer)列为反模式。
Project Reactor 的 Flux/Mono 实现 Publisher;subscribe() 返回的 Disposable 封装 Subscription.cancel()。StepVerifier 测试即验证协议合规,而非仅断言最终值。
// 规范方法在 Reactor 高层的体现 Disposable d = Flux.range(1, 100) .take(10) // 等价于在规范层只 request 10 个 .subscribe(System.out::println); // 返回 Disposable = Subscription 包装 d.dispose(); // 等价于 subscription.cancel()
Reacton 没有让你手写 request/cancel,而是把规范语义折叠进算子:take(n) 自动按需请求,limitRate(n) 自动分批请求,subscribe() 返回的 Disposable 就是 Subscription 的便捷包装。这套设计的价值在于:你可以写「符合规范」的代码而不接触规范——但调试时你仍要知道底层发生了什么,比如为什么 take(10) 之后上游不再收到 request,为什么 dispose() 之后源端的连接被关闭。规范是地板,框架是装修,你踩在地板上,但要知道地板在哪。

| 检查项 | 违反时现象 | 测试手段 |
|---|---|---|
| 串行调用 | 并发 onNext,值乱序 | TCK / StepVerifier |
| request 不超配额 | Publisher 超发 | 计数断言 |
| 错误后不 onNext | 错误后又来值 | expectError 后无 expectNext |
| cancel 后停发 | 泄漏、资源不释放 | thenCancel() + 资源断言 |
Reactive Streams 提供官方 TCK(Technology Compatibility Kit),任何声称兼容的库都要跑这套一致性测试。Reacton 的 StepVerifier 在功能测试层面覆盖同样的断言维度——这正是它被称为「验证协议合规」的原因:你测的不是「值对不对」,而是「信号的顺序、次数、终止是否符合契约」。
为了便于记忆与自测,把四个接口的 13 个方法完整列出:
| 接口 | 方法 | 数量 |
|---|---|---|
Publisher<T> |
subscribe(Subscriber) |
1 |
Subscriber<T> |
onSubscribe(Subscription)、onNext(T)、onError(Throwable)、onComplete() |
4 |
Subscription |
request(long)、cancel() |
2 |
Processor<T,R> |
继承 Publisher + Subscriber | 6(合计) |
合计 1 + 4 + 2 + 6 = 13。闭卷回忆这 13 个方法,是最快的规范熟悉度测试——如果你能默写出来并说出每个方法的触发时机与约束,你已经超过大多数「会用响应式」的开发者的水平。
// 规范中最重要的几条红线 // 1. request(n) 的 n 必须为正数 subscription.request(0); // 违规:n <= 0 是非法请求 // 2. 信号串行:onNext 之间不得并发 // 3. onError/onComplete 后不得再有 onNext // 4. cancel() 后 Publisher 必须停发 subscription.cancel(); // 此后不得再调用 subscription 的任何方法
「串行调用」是规范最容易忽略却最重要的条款:它保证 Subscriber 看到的信号是线性排列的,从而允许开发者用单线程心智推理整条链。实现上这意味着 Publisher 必须用队列、锁或线程限制来串行化来自多线程源的信号——这也是自定义 Publisher 最常出错的点。
onError 是终结判决,影响故障隔离与重试设计一句话总结:Reactive Streams 用 4 接口 13 方法,把「发布/订阅如何对话」从各框架的私约变成整个 JVM 生态的公法。
下一节看 R2DBC、MQTT 等如何把同一语义扩展到数据库与物联网。