Project Reactor是Java响应式编程库,提供Mono和Flux核心类型,支持非阻塞、背压及异步数据流处理。它是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 体系中最基础、也最容易被低估的一节——它看似简单,却蕴含了响应式编程最核心的两条法则:
- 声明 ≠ 执行:创建 Flux/Mono 只是画蓝图,不会触发任何计算。
- 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──|> onComplete3.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 → Mono3.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/Set | Flux.fromIterable(collection) | 遍历集合 |
| 整数序列 | Flux.range(start, count) | 生成区间 |
| 可能为 null | Mono.justOrEmpty(value) | 安全包装 |
| 惰性计算 | Mono.fromSupplier(() -> …) | 每次订阅重新计算 |
| 阻塞调用 | Mono.fromCallable(() -> …) | 配合 boundedElastic |
| CompletableFuture | Mono.fromFuture(future) | 桥接异步 API |
| 回调式 API | Mono.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 的第一课,也是最重要的一课。它建立了一个核心心智模型:
响应式编程 = 声明式地描述数据流 + 在正确的时机触发执行。
掌握了这个模型,后续的操作符组合、错误处理、背压控制、调度器切换,都不过是在这张蓝图上添加更精细的构件罢了。