news 2026/8/18 12:32:09

创建 Flux/Mono 并订阅:Project Reactor 响应式编程的第一步

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
创建 Flux/Mono 并订阅:Project Reactor 响应式编程的第一步

Project Reactor是Java响应式编程库,提供MonoFlux核心类型,支持非阻塞、背压及异步数据流处理。它是Spring WebFlux的基础,适用于构建高并发低延迟的微服务与事件驱动应用,遵循Reactive Streams规范。

1. 响应式编程入门:从阻塞困境到数据流之美
2. 深入 Project Reactor:从原理到工程实践的全面指南
3. Flux 与 Mono:Project Reactor 核心响应式类型深度解析
4. Mono:Project Reactor 中最精巧的响应式原语
5. 创建 Flux/Mono 并订阅:Project Reactor 响应式编程的第一步
6. 程序化创建响应式序列:Flux.generate、Flux.create 与 Flux.push 深度解析
7. 线程调度与 Schedulers:Project Reactor 并发模型的核心引擎
8. 响应式流中的错误处理:Project Reactor 异常治理全体系
9. Sinks API:Project Reactor 中程序化发射数据的现代方案

引言

每一位响应式编程的旅程,都始于同一个问题:

“我如何创建一个 Flux 或 Mono,然后让它真正跑起来?”

Project Reactor 官方文档在 Core Features 章节中,以 “Simple ways to create a Flux or Mono and subscribe to it” 为题,用最精炼的示例回答了这个问题。这是整个 Reactor 体系中最基础、也最容易被低估的一节——它看似简单,却蕴含了响应式编程最核心的两条法则:

  1. 声明 ≠ 执行:创建 Flux/Mono 只是画蓝图,不会触发任何计算。
  2. Subscribe 是引擎的点火钥匙:只有订阅,数据才会流动。

本文将围绕官方文档的核心内容,系统梳理 Flux 和 Mono 的创建方式、subscribe 的多种重载形式,以及背后隐藏的响应式执行模型。

一、核心概念:Nothing Happens Until You Subscribe

在动手写代码之前,必须先建立一个关键认知:

// 这行代码执行后,什么都不会发生。没有输出,没有副作用。Flux<Integer>ints=Flux.just(1,2,3);

Reactor 中的 Flux 和 Mono 是惰性的(Lazy)。它们是对数据流的声明式描述,而非数据本身。你可以把它理解为:

  • Flux/Mono = 一份"施工图纸"
  • subscribe() = “开工指令”

在没有 subscribe() 之前,操作符链只是一段内存中的对象图,不消耗任何 I/O、不占用线程、不产生输出。

二、创建 Flux 的简单方式

2.1 Flux.just():从已知元素创建

Flux<Integer>ints=Flux.just(1,2,3);

这是最直观的方式——将若干已知元素包装为一个 Flux。该 Flux 会依次发射 1 → 2 → 3,然后发出 onComplete 信号。
Marble 图:

──1──2──3──|> onComplete

💡 Flux.just() 接受可变参数,支持 1 到 N 个元素。

2.2 Flux.fromIterable():从集合创建

Flux<Integer>ints=Flux.fromIterable(Arrays.asList(1,2,3));

适用于已有 Iterable(如 List、Set)数据的场景。注意:Flux 会遍历该集合,但不会修改它。

2.3 Flux.range():从整数区间创建

Flux<Integer>ints=Flux.range(1,3);// 发射 1, 2, 3

第一个参数是起始值,第二个参数是元素个数(不是结束值)。

2.4 其他常用创建方式速览

// 空 Flux(立即 onComplete,无元素)Flux<String>empty=Flux.empty();// 错误 Flux(立即 onError)Flux<Object>error=Flux.error(newRuntimeException("boom"));// 从数组创建Flux<String>fromArray=Flux.fromArray(newString[]{"a","b","c"});// 从 Java Stream 创建Flux<String>fromStream=Flux.fromStream(list.stream());// 定时发射(无限流)Flux<Long>ticks=Flux.interval(Duration.ofSeconds(1));

三、创建 Mono 的简单方式

3.1 Mono.just():包装单个值

Mono<String>mono=Mono.just("foo");

发射一个元素 “foo”,然后 onComplete。
Marble 图:

──foo──|> onComplete

3.2 Mono.empty():空 Mono

Mono<String>empty=Mono.empty();

不发射任何元素,直接 onComplete。语义上等价于 Optional.empty() 的异步版本。

3.3 Mono.error():错误 Mono

Mono<Object>error=Mono.error(newIllegalArgumentException("invalid input"));

不发射元素,直接发出 onError 信号。

3.4 Mono.justOrEmpty():安全包装可能为 null 的值

Mono<String>safe=Mono.justOrEmpty(nullableValue);// null → Mono.empty()Mono<String>safeOpt=Mono.justOrEmpty(Optional.of("hi"));// Optional → Mono

3.5 Mono.fromSupplier():惰性计算

Mono<Long>lazy=Mono.fromSupplier(()->System.currentTimeMillis());// 每次 subscribe 时才调用 supplier

四、Subscribe:让数据真正流动

4.1 最简订阅:消费数据

官方文档给出的第一个 subscribe 示例:

Flux<Integer>ints=Flux.just(1,2,3);ints.subscribe(System.out::println);

输出:
subscribe(Consumer) 是最简形式——只处理 onNext 信号,忽略完成和错误。

4.2 订阅:数据 + 错误处理

Flux<Integer>ints=Flux.just(1,2,0,4);ints.map(i->"100 / "+i+" = "+(100/i)).subscribe(System.out::println,// onNext: 消费数据error->System.err.println("Error: "+error)// onError: 处理异常);

输出:

100 / 1 = 100 100 / 2 = 50 Error: java.lang.ArithmeticException: / by zero

当 i = 0 时触发除零异常,流终止,onError 被调用。

4.3 订阅:数据 + 错误 + 完成

Flux.just(1,2,3).subscribe(data->System.out.println("Next: "+data),// onNexterror->System.err.println("Error: "+error),// onError()->System.out.println("Done!")// onComplete);

输出:

Next: 1 Next: 2 Next: 3 Done!

4.4 订阅:完整控制(含 Subscription)

Flux.just(1,2,3,4,5).subscribe(data->System.out.println("Next: "+data),error->System.err.println("Error: "+error),()->System.out.println("Done!"),subscription->{System.out.println("Subscribed!");subscription.request(2);// 背压:只请求 2 个元素});

第四个参数是 Subscription 的回调,允许你控制初始请求量——这是背压机制的入口。

4.5 使用 BaseSubscriber 精细控制

Flux.just(1,2,3,4,5).subscribe(newBaseSubscriber<Integer>(){@OverrideprotectedvoidhookOnSubscribe(Subscriptionsubscription){System.out.println("Subscribed");request(2);// 初始请求 2 个}@OverrideprotectedvoidhookOnNext(Integervalue){System.out.println("Received: "+value);request(1);// 每处理完一个,再要一个}@OverrideprotectedvoidhookOnComplete(){System.out.println("All done");}@OverrideprotectedvoidhookOnError(Throwablethrowable){System.err.println("Error: "+throwable);}});

五、subscribe 重载形式全景图

签名用途
subscribe()触发订阅,不处理任何信号(极少使用)
subscribe(Consumer)只消费 onNext
subscribe(Consumer, Consumer)消费数据 + 处理错误
subscribe(Consumer, Consumer, Runnable)数据 + 错误 + 完成
subscribe(Consumer, Consumer, Runnable, Consumer)全量控制,含初始 request
subscribe(Subscriber)传入完整 Subscriber 实现

六、完整示例:从创建到订阅的端到端流程

6.1 Flux 示例

importreactor.core.publisher.Flux;publicclassFluxExample{publicstaticvoidmain(String[]args){// 1. 创建(声明)Flux<Integer>ints=Flux.just(1,2,3);System.out.println("Flux created, but nothing happened yet.");// 2. 添加操作符(仍然是声明)Flux<String>transformed=ints.map(i->"Value: "+i);System.out.println("Operators added, still nothing happened.");// 3. 订阅(执行!)System.out.println("--- Subscribing now ---");transformed.subscribe(System.out::println,error->System.err.println("Error: "+error),()->System.out.println("=== Complete ==="));}}

输出:

Flux created, but nothing happened yet. Operators added, still nothing happened. --- Subscribing now --- Value: 1 Value: 2 Value: 3 === Complete ===

6.2 Mono 示例

importreactor.core.publisher.Mono;publicclassMonoExample{publicstaticvoidmain(String[]args){// 创建Mono<String>mono=Mono.just("Hello, Reactor!");// 订阅mono.subscribe(value->System.out.println("Got: "+value),error->System.err.println("Failed: "+error),()->System.out.println("Mono completed."));}}

输出:

Got: Hello, Reactor! Mono completed.

6.3 错误传播示例

Flux.just(10,5,0,2).map(i->100/i).subscribe(result->System.out.println("Result: "+result),error->System.out.println("Caught: "+error.getMessage()),()->System.out.println("Done"));

输出:

Result: 10 Result: 20 Caught: / by zero

注意:错误发生后,流立即终止。Done 不会被打印,onComplete 和 onError 互斥。

七、subscribe 的内部执行机制

当 subscribe() 被调用时,Reactor 内部执行以下步骤:

┌─────────────────────────────────────────────────────────────────┐ │ subscribe() 触发后的流程 │ ├─────────────────────────────────────────────────────────────────┤ │ │ │ 1. 构建操作符链(如果尚未构建) │ │ source → operator1 → operator2 → ... → terminal │ │ │ │ 2. 从终端向上游逐层调用 subscribe() │ │ terminal.subscribe() → op2.subscribe() → op1.subscribe() │ │ → source.subscribe() │ │ │ │ 3. 源头 Publisher 调用 Subscriber.onSubscribe(Subscription) │ │ → 传递 Subscription 对象给下游 │ │ │ │ 4. 下游通过 Subscription.request(n) 请求数据 │ │ → 默认 subscribe(Consumer) 请求 Long.MAX_VALUE(无界) │ │ │ │ 5. 数据沿链从上游流向下游:onNext(T) │ │ source → op1.transform → op2.transform → subscriber.onNext │ │ │ │ 6. 终止信号:onComplete 或 onError │ │ → 清理资源,流生命周期结束 │ │ │ └─────────────────────────────────────────────────────────────────┘

关键点:

  • 订阅信号(subscribe)是自下游向上游传播的
  • 数据信号(onNext)是自上游向下游流动的
  • 这是一个双向握手过程

八、Disposable:订阅的生命周期管理

subscribe() 方法返回一个 Disposable 对象,用于取消订阅:

Flux<Long>infiniteStream=Flux.interval(Duration.ofMillis(100));Disposabledisposable=infiniteStream.subscribe(tick->System.out.println("Tick: "+tick));// 500ms 后取消Thread.sleep(500);disposable.dispose();// 取消订阅,停止数据流System.out.println("Disposed. Stream stopped.");

⚠️ 对于无限流(如 Flux.interval()),必须保存 Disposable 并在适当时机调用 dispose(),否则流将永远运行,造成资源泄漏。

Disposable 的常用方法:

方法说明
dispose()取消订阅,触发 onCancel 信号向上游传播
isDisposed()查询是否已取消

九、常见错误与注意事项

9.1 ❌ 忘记 subscribe

// 这段代码不会有任何输出!Flux.just(1,2,3).map(i->i*10);// 没有 subscribe → 什么都不会发生

9.2 ❌ 多次 subscribe 导致重复执行

Flux<String>flux=Flux.fromCallable(()->{System.out.println("Executing expensive operation...");returnfetchFromDatabase();});flux.subscribe(System.out::println);// 第一次执行flux.subscribe(System.out::println);// 第二次执行(Cold Publisher)// "Executing expensive operation..." 会打印两次!

解决方案:使用 .cache() 或 .share() 将 Cold 转为 Hot。

9.3 ❌ 在 subscribe 中抛出未捕获异常

Flux.just(1,2,3).subscribe(i->{if(i==2)thrownewRuntimeException("unexpected");System.out.println(i);});// 异常会被 Reactor 的 onErrorDropped 钩子捕获,可能导致意外行为

**正确做法:**在操作符链中用 onErrorResume / onErrorReturn 处理。

9.4 ❌ 在响应式链中使用 block()

// ❌ 在 WebFlux / Netty EventLoop 中绝对禁止Mono<String>result=someMono.block();// 阻塞线程!

十、从"简单创建"到工程实践

10.1 Spring WebFlux 中的自动订阅

在 WebFlux Controller 中,你不需要手动 subscribe——框架会自动完成:

@GetMapping("/users/{id}")publicMono<User>getUser(@PathVariableLongid){returnuserRepository.findById(id);// 返回 Mono,框架负责 subscribe}

10.2 测试中的 StepVerifier

@TestvoidtestFluxCreation(){StepVerifier.create(Flux.just(1,2,3)).expectNext(1).expectNext(2).expectNext(3).verifyComplete();}@TestvoidtestMonoCreation(){StepVerifier.create(Mono.just("hello")).expectNext("hello").verifyComplete();}@TestvoidtestEmptyMono(){StepVerifier.create(Mono.empty()).verifyComplete();// 期望直接完成,无元素}@TestvoidtestErrorFlux(){StepVerifier.create(Flux.error(newRuntimeException("oops"))).expectErrorMessage("oops").verify();}

10.3 桥接阻塞代码的正确姿势

// 将阻塞调用包装为 Mono,并切换到弹性线程池Mono<String>result=Mono.fromCallable(()->{// 阻塞操作:JDBC 查询、文件读取、HTTP 同步调用returnlegacyService.blockingFetch();}).subscribeOn(Schedulers.boundedElastic());result.subscribe(System.out::println);

十一、创建方式选择指南

场景推荐方式说明
已知固定元素Flux.just(a, b, c)最直接
已有 List/SetFlux.fromIterable(collection)遍历集合
整数序列Flux.range(start, count)生成区间
可能为 nullMono.justOrEmpty(value)安全包装
惰性计算Mono.fromSupplier(() -> …)每次订阅重新计算
阻塞调用Mono.fromCallable(() -> …)配合 boundedElastic
CompletableFutureMono.fromFuture(future)桥接异步 API
回调式 APIMono.create(sink -> …)桥接第三方 SDK
空结果Mono.empty() / Flux.empty()无数据但成功
立即失败Mono.error(e) / Flux.error(e)校验失败等
动态决策Mono.defer(() -> …)每次订阅时选择不同源

十二、总结

┌──────────────────────────────────────────────────────────────────────┐ │ │ │ 创建 Flux/Mono → 添加操作符 → subscribe() │ │ (声明蓝图) (描述变换) (点火执行) │ │ │ │ Flux.just(1,2,3) .map(...) .subscribe( │ │ Mono.just("hi") .filter(...) onNext, │ │ Flux.fromIterable(list) .flatMap(...) onError, │ │ Mono.fromSupplier(s) .onErrorResume(...) onComplete │ │ Mono.defer(...) ) │ │ │ │ ⚠️ 没有 subscribe → 什么都不会发生 │ │ ⚠️ subscribe 返回 Disposable → 管理生命周期 │ │ ⚠️ 在 WebFlux 中 → 框架自动 subscribe,不要手动调用 │ │ ⚠️ 不要 block() → 保持全链路非阻塞 │ │ │ └──────────────────────────────────────────────────────────────────────┘

"创建 Flux/Mono 并订阅"是 Project Reactor 的第一课,也是最重要的一课。它建立了一个核心心智模型:

响应式编程 = 声明式地描述数据流 + 在正确的时机触发执行。

掌握了这个模型,后续的操作符组合、错误处理、背压控制、调度器切换,都不过是在这张蓝图上添加更精细的构件罢了。

版权声明: 本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若内容造成侵权/违法违规/事实不符,请联系邮箱:809451989@qq.com进行投诉反馈,一经查实,立即删除!
网站建设 2026/8/18 12:29:21

非线性层扩散:多智能体系统辨识的新范式与工程实践

1. 项目概述&#xff1a;当多智能体系统遇上非线性层扩散 最近在复现和思考一些前沿的分布式系统辨识工作时&#xff0c;一个绕不开的概念就是“非线性层扩散”。乍一听&#xff0c;这像是把“非线性系统”和“层理论”这两个硬核概念强行捏合在一起&#xff0c;充满了数学的“…

作者头像 李华
网站建设 2026/8/18 12:20:38

C++26新特性落地,存量代码库的现代升级经验

从标准到产线&#xff1a;C26 新特性的工程化落地思路 C26 的草案已趋于稳定&#xff0c;对于长期维护百万行级代码库的系统软件团队来说&#xff0c;真正的挑战从来不是“特性有没有”&#xff0c;而是“怎么在不翻车的前提下用上新东西”。过去两年&#xff0c;我所在的团队…

作者头像 李华
网站建设 2026/8/18 12:18:02

Packmol 从零上手:分子动力学模拟初始构型的完整构建指南

Packmol 从零上手&#xff1a;分子动力学模拟初始构型的完整构建指南 【免费下载链接】packmol Packmol - Initial configurations for molecular dynamics simulations 项目地址: https://gitcode.com/gh_mirrors/pa/packmol Packmol 是一个用 Fortran 写就的开源打包工…

作者头像 李华
网站建设 2026/8/18 12:17:57

2026上海移动硬盘数据恢复选型避坑指南及机构对比测试

2026年&#xff0c;高频跨城办公与大型项目协作让便携式移动硬盘成为职场刚需标配。伴随频繁携带&#xff0c;摔落异响、意外进水等物理损坏&#xff0c;以及误格式化引起的逻辑故障成为数据安全的高危触点。当意外发生时&#xff0c;一块存储着核心项目代码或绝密标书的硬盘&a…

作者头像 李华
网站建设 2026/8/18 12:17:11

大模型服务新范式:Agentic Discovery实现推理时动态协同与性能扩展

1. 项目概述&#xff1a;当大模型开始自我进化最近在跟几个做模型推理和部署的朋友聊天&#xff0c;大家普遍有个感觉&#xff1a;大模型&#xff08;LLMs&#xff09;上线后的表现&#xff0c;跟离线评测时总有点“货不对板”。你精心调教好的模型&#xff0c;在真实用户五花八…

作者头像 李华