本节摘要:Project Reactor 用
Mono<T>(0–1 元素)与Flux<T>(0–N 元素)构成流代数;每个流的生命周期由 onSubscribe、onNext、onError、onComplete 信号描述。本节建立 Publisher/Subscriber 直觉,为背压与规范章打地基。
Reactive Streams 四角色中最核心是:
onComplete 或 onError 终止Observable 不是静态容器,而是封装了订阅契约、错误路径与生命周期的管道。
| 类型 | 元素个数 | 典型场景 |
|---|---|---|
Mono<T> |
0 或 1 | 单次 HTTP、DB 单行 |
Flux<T> |
0..∞ | SSE、Kafka 消费、传感器 |
Reactor 将 map、filter、flatMap、zip 设计为满足结合律的纯函数操作,形成可组合、可测试的流代数(SOURCE 第二章)。
Flux.range(1, 5) .map(i -> i * 2) .filter(i -> i > 4) .subscribe(System.out::println); // 输出: 6, 8, 10
「Mono 还是 Flux」这个选择不是风格问题,而是契约问题。当你把方法签名写成 Mono<User>,你同时向调用方承诺三件事:结果可能没有(empty)、最多一个、错误时会以 onError 信号结束。写成 Flux<User> 则承诺:可能有很多个、按需请求(request(n))、中间可被取消。调用方依据这个签名就能决定怎么处理结果——这是类型驱动的 API 设计,和「返回一个 List 还是 Optional」的思考同构但更强,因为信号还带着时间。
一个订阅关系建立后,Subscriber 收到的信号只有四种,它们的组合构成流生命周期的全部:
| 信号 | 携带内容 | 次数约束 | 含义 |
|---|---|---|---|
onSubscribe |
Subscription |
恰好 1 次 | 订阅建立,获得 request/cancel 能力 |
onNext |
元素 T | 0..N 次 | 每来一个值 |
onError |
Throwable |
最多 1 次(与 onComplete 互斥) | 失败终止 |
onComplete |
无 | 最多 1 次 | 正常终止 |
信号顺序保证:onNext 之间不会并发,onError/onComplete 之后不再有 onNext。这条「串行保证」让开发者在单线程心智下思考整个链路,不用处理回调里的竞态——这是响应式框架替你管理并发的重要部分。
// 手工 Subscriber:体会四个回调 Flux.just("a", "b") .subscribe(new BaseSubscriber<>() { @Override protected void hookOnSubscribe(Subscription s) { System.out.println("已订阅,请求 1 个"); s.request(1); } @Override protected void hookOnNext(String v) { System.out.println("收到: " + v); request(1); // 逐个请求,天然背压 } @Override protected void hookOnComplete() { System.out.println("流结束"); } });
这个 BaseSubscriber 逐条对应信号模型:hookOnSubscribe 拿到 Subscription 并 request(1);hookOnNext 消费一个元素后再次 request(1)——注意,这就是背压的最小实现,消费者用「每次请求一个」主动控制节奏;hookOnComplete 处理正常结束。整个过程中没有回调嵌套、没有共享状态,四个钩子的执行顺序由框架保证串行。如果你能完全看懂这段代码,2.3 的背压对你就不再神秘。
| 概念 | 命令式 | 响应式 |
|---|---|---|
| 数据到达 | 返回值 | onNext 信号 |
| 结束 | return | onComplete |
| 失败 | throw | onError(终止性) |
| 流量 | 无 | request(n) |
⚠️ 常见坑:链末端
subscribe()不传onError,异常可能进全局钩子而难排查。
💡 关键直觉:
Flux类型签名本身声明了「可能很多值、可能很长生命周期」——影响后续能否随便cache()或share()(见第 5 章)。

初学者最容易犯的思维错误,是把 Flux 当成「异步 List」——然后就奇怪为什么遍历一次就没元素了。区分两者:
判断口诀:「这个数据源现在就有全部数据吗?」 是 → 容器;「数据要随时间/事件产生吗?」 是 → 流。Flux.just(1,2,3) 虽然元素是现成的,但依然承诺了订阅触发语义,所以仍是流。这条区分直接决定你该不该 cache()——容器天然可反复读,流必须显式缓存才能反复订阅。
理解了四种信号,就理解了为什么响应式算子可以任意组合。每个算子本质上是一个「信号变换器」:它订阅上游的四种信号,经过变换后向下游发出自己的四种信号。例如 map 就是把上游的每个 onNext(x) 换成 onNext(f(x)),onError/onComplete 原样透传。因为所有算子的输入输出都是同一套信号协议,所以算子之间可以任意拼接——这就是「流代数」的工程基础。
// 算子即信号变换器:filter 中断某些 onNext Flux.just(1, 2, 3, 4, 5) .filter(i -> i % 2 == 0) // onNext(1,3,5) 被过滤掉 .map(i -> i * 10) // onNext(2) → onNext(20) ... .subscribe(System.out::println); // 输出: 20, 40
| 概念 | Reactor | RxJS | 语义 |
|---|---|---|---|
| 发布者 | Mono/Flux |
Observable |
惰性、订阅才执行 |
| 订阅 | subscribe() |
subscribe() |
触发执行 |
| 元素信号 | onNext |
next |
携带数据 |
| 完成信号 | onComplete |
complete |
正常终止 |
| 错误信号 | onError |
error |
终止 + 携带异常 |
| 取消 | dispose() |
unsubscribe() |
停止上游 |
这张对齐表能帮你在前后端之间自如切换:同一个「搜索防抖」需求,后端用 Reactor 写 switchMap,前端用 RxJS 写 switchMap,信号模型完全同构——差异只在背压与调度(第 4 章详述)。这也是为什么理解信号模型比记住某个框架的 API 更重要:信号模型是响应式的世界语,框架只是方言。
下一节讨论 Scheduler 如何仲裁物理时间与测试时间。