【WebFlux】第二篇 —— Project Reactor 核心数据类型与doOnXXX介绍

【WebFlux】第二篇 —— Project Reactor 核心数据类型与doOnXXX介绍
认识 Project Reactor响应式流的“基石”在 Spring WebFlux 的底层真正支撑起异步非阻塞数据流转的是一个名为 Project Reactor 的核心库。它完全实现了 Reactive Streams 规范为我们提供了一套声明式、函数式的 API。如果说 Reactive Streams 是响应式编程的“交通规则”那么 Project Reactor 就是按照这套规则制造出的“超级跑车”。在 Reactor 中万物皆流。为了应对不同的数据场景Reactor 提供了两个最核心的数据类型PublisherMono和Flux。它们是整个响应式编程大厦的基石。核心类型解析Mono 与 Flux要掌握 Reactor首先要分清这两个核心概念的区别Mono0 或 1 个元素的异步序列Mono 代表一个最多只包含单个元素的异步计算结果。你可以把它理解为异步版的 Optional 或 CompletableFuture。典型场景根据 ID 查询单个用户信息、保存一条记录、执行一次无返回值的异步操作如 Mono、HTTP 接口返回单个对象等。Flux0 到 N 个元素的异步序列Flux 代表一个包含 0 到多个元素的有序异步序列它甚至可以是一个无限流。你可以把它想象成一条物流传送带或者数据库的游标。典型场景查询用户列表、处理文件中的多行数据、WebSocket 消息流、实时传感器数据推送等。声明式与惰性执行Lazy Evaluation这是响应式编程中最反直觉、但也最核心的特性。在 Reactor 中当你调用 map、filter 等操作符时实际上并没有任何数据被处理也没有任何业务逻辑被执行。这些操作仅仅是在构建一条“处理流水线Pipeline”。只有当有 Subscriber订阅者调用 subscribe() 方法时整条流水线才会被激活数据才会像水流一样从源头开始向下流动。这种惰性执行机制使得我们可以像搭积木一样灵活地组装和复用数据流逻辑同时也避免了不必要的资源消耗。数据流的生命周期与“弹珠图Marble Diagrams”doOnXXX是响应式流里的副作用side-effect观察者它监听信号经过、不修改流、不改变元素返回的是同一个流。用途打日志、埋点、调试、资源清理。铁律改造用map/filter观察才用doOnXXX不subscribe不触发冷流。速查表收藏级方法触发信号典型用途doFirst订阅前仅 1 次一次性前置准备doOnSubscribe订阅初始化、拿 SubscriptiondoOnRequest下游请求观察背压doOnCancel取消清理doOnNext每个元素日志 / 埋点doOnEach所有 Signal全信号观察调试doOnComplete正常结束收尾doOnError错误错误日志降级前可见doOnTerminate终止前终止前打扫doAfterTerminate终止后终止后打扫doOnSuccessMono 成功Mono 收尾值可能为 nulldoFinally任意终止兜底清理带SignalTypedoOnDiscard元素被丢弃释放被丢元素持有的资源按信号阶段分类13 个方法阶段方法订阅 / 生命周期doFirst、doOnSubscribe、doOnRequest、doOnCancel元素级doOnNext、doOnEach终止doOnComplete、doOnError、doOnTerminate、doAfterTerminate、doOnSuccess(Mono)、doFinally资源清理doOnDiscard复合案例Case A · 正常完成的全生命周期一次演示 8 个方法Flux.range(1,3).doFirst(()-System.out.println([doFirst] 订阅前执行一次)).doOnSubscribe(s-System.out.println([doOnSubscribe] s)).doOnRequest(n-System.out.println([doOnRequest] 请求了 n)).doOnNext(i-System.out.println([doOnNext] 元素 i)).doOnComplete(()-System.out.println([doOnComplete] 正常结束)).doOnTerminate(()-System.out.println([doOnTerminate] 即将终止)).doAfterTerminate(()-System.out.println([doAfterTerminate] 已下发)).doFinally(type-System.out.println([doFinally] 类型type)).subscribe(v-System.out.println( 消费 v));典型输出[doFirst] 订阅前执行一次 [doOnSubscribe] reactor.core.publisher.FluxRange$RangeSubscription... [doOnRequest] 请求了 9223372036854775807 [doOnNext] 元素 1 消费 1 [doOnNext] 元素 2 消费 2 [doOnNext] 元素 3 消费 3 [doOnComplete] 正常结束 [doOnTerminate] 即将终止 [doAfterTerminate] 已下发 [doFinally] 类型ON_COMPLETEdoOnRequest的9223372036854775807Long.MAX_VALUE即默认subscribe是无限请求背压全开。顺序口诀先订阅 → 后请求 → 逐个 next/消费 → 完成前 terminate → 下发后 afterTerminate → 最后 finally。Case B · 错误路径全家桶一次演示 5 个方法Flux.just(2,0).map(i-10/i)// i0 时抛 ArithmeticException.doOnNext(i-System.out.println([doOnNext] i)).doOnError(e-System.out.println([doOnError] e.getMessage())).doOnTerminate(()-System.out.println([doOnTerminate] 即将终止)).doAfterTerminate(()-System.out.println([doAfterTerminate] 已下发)).doFinally(type-System.out.println([doFinally] 类型type)).onErrorResume(e-Flux.just(-1))// 降级.subscribe(v-System.out.println( 消费 v));典型输出[doOnNext] 5 消费 5 [doOnError] / by zero [doOnTerminate] 即将终止 [doAfterTerminate] 已下发 [doFinally] 类型ON_ERROR 消费 -1doOnError在onErrorResume降级之前就看到了原始异常。doFinally的ON_ERROR也先于降级恢复触发——它告诉你这条流是怎么死的。Case C · Mono 成功与收尾一次演示 4 个方法Mono.just(hello).doOnSubscribe(s-System.out.println([doOnSubscribe])).doOnNext(v-System.out.println([doOnNext] v)).doOnSuccess(v-System.out.println([doOnSuccess] 成功值v)).doFinally(type-System.out.println([doFinally] 类型type)).subscribe();Mono.empty().doOnSuccess(v-System.out.println([doOnSuccess] 空完成vv))// v null.subscribe();输出非空 Mono[doOnSubscribe] [doOnNext] hello [doOnSuccess] 成功值hello [doFinally] 类型ON_COMPLETE输出空 Mono[doOnSuccess] 空完成vnulldoOnSuccessdoOnNextdoOnComplete的合体Mono 专用。空完成时v null判空要小心。Case D · 背压与丢弃doOnRequestdoOnDiscardFlux.range(1,100).doOnRequest(n-System.out.println([doOnRequest] n)).doOnDiscard(Integer.class,i-System.out.println([doOnDiscard] 丢弃 i)).onBackpressureDrop()// 下游不取就丢弃.subscribe(newBaseSubscriberInteger(){OverrideprotectedvoidhookOnSubscribe(Subscriptions){s.request(2);}OverrideprotectedvoidhookOnNext(Integerv){System.out.println( 消费 v);}});输出下游只取 2 个其余被丢弃[doOnRequest] 2 消费 1 消费 2 [doOnDiscard] 丢弃 3 [doOnDiscard] 丢弃 4 ...5~100 同理被丢弃doOnDiscard在元素因背压丢弃或取消时触发用来释放元素持有的资源连接、句柄避免泄漏。Case E · 取消doOnCancelDisposabledFlux.interval(Duration.ofMillis(100)).doOnCancel(()-System.out.println([doOnCancel] 被取消了)).subscribe(v-System.out.println( v));Thread.sleep(350);d.dispose();// 主动取消 → 触发 doOnCancelCase F · 全信号监听一个doOnEach顶全部信号类方法Flux.just(1,2,3).doOnEach(signal-{switch(signal.getType()){caseON_NEXT:System.out.println(NEXT signal.get());break;caseON_COMPLETE:System.out.println(COMPLETE);break;caseON_ERROR:System.out.println(ERROR signal.getThrowable().getMessage());break;default:System.out.println(其他 signal.getType());}}).subscribe();doOnEach(Signal)把订阅、请求、取消、next、complete、error 全部包成Signal对象。调试全信号时直接上log()内置的完整信号日志或doOnEach即可不必逐个手写。三大必踩的坑坑 1线程上下文随publishOn位置变化Flux.range(1,2).doOnNext(i-System.out.println(A 线程Thread.currentThread().getName() 值i)).publishOn(Schedulers.parallel()).doOnNext(i-System.out.println(B 线程Thread.currentThread().getName() 值i)).blockLast();输出A在订阅线程跑B在 parallel 线程跑。doOnNext在链上的位置决定它在哪条线程执行——排查日志重复/顺序乱时这是第一怀疑点。坑 2doOnXXX内抛异常会污染整条流Flux.just(1,2).doOnNext(i-{if(i2)thrownewRuntimeException(炸了);})// 会让流直接 error.subscribe(v-{},e-System.out.println(收到错误: e));doOnNext里抛异常会变成 error 信号向上游传播整个序列挂掉。所以doOnXXX里只放轻量、不会失败的逻辑。坑 3冷流不订阅不触发上面所有doOnXXX都只在.subscribe()后才执行。组装好链式但忘了订阅 什么都不会发生。丰富的数据源创建方式Reactor 提供了极其丰富的工厂方法来创建 Mono 和 Flux以适配各种业务场景静态值创建使用 Mono.just(“Hello”) 或 Flux.just(“A”, “B”, “C”) 包装已知数据。// 创建包含单个元素的 MonoMono.just(Hello WebFlux).subscribe(System.out::println);// 创建包含多个元素的 FluxFlux.just(Java,Go,Rust).subscribe(System.out::println);空流与错误流使用 Mono.empty() 表示无数据返回使用 Mono.error(new RuntimeException()) 直接抛出异常信号。// 创建空流订阅后直接触发 onCompleteMono.empty().subscribe(data-{},error-{},()-System.out.println(空流已完成));// 创建错误流订阅后直接触发 onErrorFlux.error(newIllegalStateException(非法状态)).subscribe();延迟/惰性初始化使用 Mono.fromSupplier(() - …) 或 Mono.defer(() - …)。这种方式只有在真正被订阅时才会执行 Supplier 内部的逻辑非常适合封装数据库查询等耗时操作。// 每次订阅都会重新执行 Supplier 中的逻辑Mono.fromSupplier(()-当前时间: System.currentTimeMillis()).subscribe(System.out::println);异步数据源转换如果系统中已有传统的异步代码可以使用 Mono.fromFuture() 或 Mono.fromCallable() 将其无缝转换为响应式流。// 包装 CompletableFutureMono.fromFuture(CompletableFuture.supplyAsync(()-异步结果)).subscribe(System.out::println);// 包装同步但耗时的 CallableMono.fromCallable(()-{Thread.sleep(1000);// 模拟耗时操作return计算完成;}).subscribe(System.out::println);时间驱动使用 Flux.interval(Duration.ofSeconds(1)) 可以创建一个每秒发射一次递增数字的无限流这在定时任务或心跳检测中非常有用。// 生成 1 到 5 的整数序列Flux.range(1,5).subscribe(i-System.out.print(i ));// 输出: 1 2 3 4 5// 每秒发射一个递增数字的无限流需配合 take 限制长度避免无限打印Flux.interval(Duration.ofSeconds(1)).take(3).subscribe(i-System.out.println(Tick: i));避坑提示警惕副作用Side Effects由于惰性执行的存在初学者极易踩坑。例如如果在 map 操作符中直接打印日志或修改外部变量这些操作只有在被订阅时才会执行。如果不小心订阅了两次这些副作用就会被执行两次。// 错误做法在 map 中执行副作用如打印日志// 问题如果该流被订阅了两次处理数据: 就会被打印两次产生不可控的副作用。Flux.just(Data-1,Data-2).map(data-{System.out.println(处理数据: data);// 副作用混入了数据转换逻辑returndata.toUpperCase();}).subscribe();// 正确做法使用 doOnNext 等生命周期钩子// 优势语义清晰doOnNext 仅作为“观察者”记录日志绝不改变流中的数据且易于在调试期移除。Flux.just(Data-1,Data-2).doOnNext(data-System.out.println(准备处理数据: data))// 安全的副作用钩子.map(String::toUpperCase)// 保持纯粹的同步转换逻辑.doOnComplete(()-System.out.println(所有数据处理完毕))// 统一处理完成事件.subscribe();最佳实践永远不要在 map 或 flatMap 中执行副作用操作。如果需要记录日志或进行调试请使用 Reactor 专门提供的“生命周期钩子”操作符如 doOnNext、doOnError、doOnComplete 等。这些钩子只会“观察”数据流而不会改变数据流本身是调试响应式代码的利器。本篇小结Mono 和 Flux 是响应式编程的容器理解了它们的惰性执行机制和生命周期信号我们就掌握了控制数据流的钥匙。下一步预告数据流建立起来了我们该如何对它们进行加工下一篇笔记我们将深入实战详解 map 与 flatMap 的核心区别并学习如何使用操作符对数据流进行转换、过滤与异常处理。