2.1 数据流抽象模型


2.1 数据流抽象模型

本节摘要:Project Reactor 用 Mono<T>(0–1 元素)与 Flux<T>(0–N 元素)构成流代数;每个流的生命周期由 onSubscribe、onNext、onError、onComplete 信号描述。本节建立 Publisher/Subscriber 直觉,为背压与规范章打地基。

先说结论

  1. 区分 Mono 与 Flux 的语义契约
  2. 写出 Subscriber 四个回调的职责
  3. 对比「容器」与「带拓扑的时空管道」两种理解

一、Publisher 与 Subscriber

Reactive Streams 四角色中最核心是:

  • Publisher:承诺发出 0..N 个元素,并以 onCompleteonError 终止
  • Subscriber:消费元素、处理终止信号、通过 Subscription 按需索取

Observable 不是静态容器,而是封装了订阅契约、错误路径与生命周期的管道

二、Reactor 类型即契约

类型 元素个数 典型场景
Mono<T> 0 或 1 单次 HTTP、DB 单行
Flux<T> 0..∞ SSE、Kafka 消费、传感器

Reactor 将 mapfilterflatMapzip 设计为满足结合律的纯函数操作,形成可组合、可测试的流代数(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 拿到 Subscriptionrequest(1)hookOnNext 消费一个元素后再次 request(1)——注意,这就是背压的最小实现,消费者用「每次请求一个」主动控制节奏;hookOnComplete 处理正常结束。整个过程中没有回调嵌套、没有共享状态,四个钩子的执行顺序由框架保证串行。如果你能完全看懂这段代码,2.3 的背压对你就不再神秘。

四、与命令式对照

概念 命令式 响应式
数据到达 返回值 onNext 信号
结束 return onComplete
失败 throw onError(终止性)
流量 request(n)

⚠️ 常见坑:链末端 subscribe() 不传 onError,异常可能进全局钩子而难排查。

💡 关键直觉Flux 类型签名本身声明了「可能很多值、可能很长生命周期」——影响后续能否随便 cache()share()(见第 5 章)。

02-02-fig01-2

容器与管道之争

初学者最容易犯的思维错误,是把 Flux 当成「异步 List」——然后就奇怪为什么遍历一次就没元素了。区分两者:

  • 容器(List/数组)是惰性的存储:元素已经存在,迭代只是读取;
  • (Flux/Observable)是惰性的行为:元素在订阅时才开始产生,迭代本身会触发副作用(发请求、读数据库、监听事件)。

判断口诀:「这个数据源现在就有全部数据吗?」 是 → 容器;「数据要随时间/事件产生吗?」 是 → 流。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 的信号对齐

概念 Reactor RxJS 语义
发布者 Mono/Flux Observable 惰性、订阅才执行
订阅 subscribe() subscribe() 触发执行
元素信号 onNext next 携带数据
完成信号 onComplete complete 正常终止
错误信号 onError error 终止 + 携带异常
取消 dispose() unsubscribe() 停止上游

这张对齐表能帮你在前后端之间自如切换:同一个「搜索防抖」需求,后端用 Reactor 写 switchMap,前端用 RxJS 写 switchMap,信号模型完全同构——差异只在背压与调度(第 4 章详述)。这也是为什么理解信号模型比记住某个框架的 API 更重要:信号模型是响应式的世界语,框架只是方言。

本章回顾

  • 信号三元组:数据、错误、完成
  • Mono/Flux 是代数类型,不是语法糖
  • 订阅建立后才发生背压协商

下一节讨论 Scheduler 如何仲裁物理时间与测试时间。


作者与出处
原作者: 灏天文库
来源:灏天文库
整理: 灏天文库整理
由灏天文库平台收录,内容或由平台用户上传,仅供学习交流
发布者: 作者: 灏天文库 转发
评论区 (0)
U