
Project Reactor是Java响应式编程库提供Mono和Flux核心类型支持非阻塞、背压及异步数据流处理。它是Spring WebFlux的基础适用于构建高并发、低延迟的微服务与事件驱动应用遵循Reactive Streams规范。1. 响应式编程入门从阻塞困境到数据流之美2. 深入 Project Reactor从原理到工程实践的全面指南3. Flux 与 MonoProject Reactor 核心响应式类型深度解析4. MonoProject Reactor 中最精巧的响应式原语5. 创建 Flux/Mono 并订阅Project Reactor 响应式编程的第一步6. 程序化创建响应式序列Flux.generate、Flux.create 与 Flux.push 深度解析7. 线程调度与 SchedulersProject Reactor 并发模型的核心引擎8. 响应式流中的错误处理Project Reactor 异常治理全体系9. Sinks APIProject 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在动手写代码之前必须先建立一个关键认知// 这行代码执行后什么都不会发生。没有输出没有副作用。FluxIntegerintsFlux.just(1,2,3);Reactor 中的 Flux 和 Mono 是惰性的Lazy。它们是对数据流的声明式描述而非数据本身。你可以把它理解为Flux/Mono 一份施工图纸subscribe() “开工指令”在没有 subscribe() 之前操作符链只是一段内存中的对象图不消耗任何 I/O、不占用线程、不产生输出。二、创建 Flux 的简单方式2.1 Flux.just()从已知元素创建FluxIntegerintsFlux.just(1,2,3);这是最直观的方式——将若干已知元素包装为一个 Flux。该 Flux 会依次发射 1 → 2 → 3然后发出 onComplete 信号。Marble 图──1──2──3──| onComplete Flux.just() 接受可变参数支持 1 到 N 个元素。2.2 Flux.fromIterable()从集合创建FluxIntegerintsFlux.fromIterable(Arrays.asList(1,2,3));适用于已有 Iterable如 List、Set数据的场景。注意Flux 会遍历该集合但不会修改它。2.3 Flux.range()从整数区间创建FluxIntegerintsFlux.range(1,3);// 发射 1, 2, 3第一个参数是起始值第二个参数是元素个数不是结束值。2.4 其他常用创建方式速览// 空 Flux立即 onComplete无元素FluxStringemptyFlux.empty();// 错误 Flux立即 onErrorFluxObjecterrorFlux.error(newRuntimeException(boom));// 从数组创建FluxStringfromArrayFlux.fromArray(newString[]{a,b,c});// 从 Java Stream 创建FluxStringfromStreamFlux.fromStream(list.stream());// 定时发射无限流FluxLongticksFlux.interval(Duration.ofSeconds(1));三、创建 Mono 的简单方式3.1 Mono.just()包装单个值MonoStringmonoMono.just(foo);发射一个元素 “foo”然后 onComplete。Marble 图──foo──| onComplete3.2 Mono.empty()空 MonoMonoStringemptyMono.empty();不发射任何元素直接 onComplete。语义上等价于 Optional.empty() 的异步版本。3.3 Mono.error()错误 MonoMonoObjecterrorMono.error(newIllegalArgumentException(invalid input));不发射元素直接发出 onError 信号。3.4 Mono.justOrEmpty()安全包装可能为 null 的值MonoStringsafeMono.justOrEmpty(nullableValue);// null → Mono.empty()MonoStringsafeOptMono.justOrEmpty(Optional.of(hi));// Optional → Mono3.5 Mono.fromSupplier()惰性计算MonoLonglazyMono.fromSupplier(()-System.currentTimeMillis());// 每次 subscribe 时才调用 supplier四、Subscribe让数据真正流动4.1 最简订阅消费数据官方文档给出的第一个 subscribe 示例FluxIntegerintsFlux.just(1,2,3);ints.subscribe(System.out::println);输出subscribe(Consumer) 是最简形式——只处理 onNext 信号忽略完成和错误。4.2 订阅数据 错误处理FluxIntegerintsFlux.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 订阅完整控制含 SubscriptionFlux.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(newBaseSubscriberInteger(){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)只消费 onNextsubscribe(Consumer, Consumer)消费数据 处理错误subscribe(Consumer, Consumer, Runnable)数据 错误 完成subscribe(Consumer, Consumer, Runnable, Consumer)全量控制含初始 requestsubscribe(Subscriber)传入完整 Subscriber 实现六、完整示例从创建到订阅的端到端流程6.1 Flux 示例importreactor.core.publisher.Flux;publicclassFluxExample{publicstaticvoidmain(String[]args){// 1. 创建声明FluxIntegerintsFlux.just(1,2,3);System.out.println(Flux created, but nothing happened yet.);// 2. 添加操作符仍然是声明FluxStringtransformedints.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){// 创建MonoStringmonoMono.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 对象用于取消订阅FluxLonginfiniteStreamFlux.interval(Duration.ofMillis(100));DisposabledisposableinfiniteStream.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 导致重复执行FluxStringfluxFlux.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(i2)thrownewRuntimeException(unexpected);System.out.println(i);});// 异常会被 Reactor 的 onErrorDropped 钩子捕获可能导致意外行为**正确做法**在操作符链中用 onErrorResume / onErrorReturn 处理。9.4 ❌ 在响应式链中使用 block()// ❌ 在 WebFlux / Netty EventLoop 中绝对禁止MonoStringresultsomeMono.block();// 阻塞线程十、从简单创建到工程实践10.1 Spring WebFlux 中的自动订阅在 WebFlux Controller 中你不需要手动 subscribe——框架会自动完成GetMapping(/users/{id})publicMonoUsergetUser(PathVariableLongid){returnuserRepository.findById(id);// 返回 Mono框架负责 subscribe}10.2 测试中的 StepVerifierTestvoidtestFluxCreation(){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并切换到弹性线程池MonoStringresultMono.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(() - …)配合 boundedElasticCompletableFutureMono.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 的第一课也是最重要的一课。它建立了一个核心心智模型响应式编程 声明式地描述数据流 在正确的时机触发执行。掌握了这个模型后续的操作符组合、错误处理、背压控制、调度器切换都不过是在这张蓝图上添加更精细的构件罢了。