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

📅 2026/7/27 8:04:52 👁️ 阅读次数 📝 编程学习
【WebFlux】第二篇 —— Project Reactor 核心数据类型与doOnXXX介绍

认识 Project Reactor:响应式流的“基石”

在 Spring WebFlux 的底层,真正支撑起异步非阻塞数据流转的,是一个名为 Project Reactor 的核心库。它完全实现了 Reactive Streams 规范,为我们提供了一套声明式、函数式的 API。如果说 Reactive Streams 是响应式编程的“交通规则”,那么 Project Reactor 就是按照这套规则制造出的“超级跑车”。

在 Reactor 中,万物皆流。为了应对不同的数据场景,Reactor 提供了两个最核心的数据类型(Publisher):MonoFlux。它们是整个响应式编程大厦的基石。

核心类型解析:Mono 与 Flux

要掌握 Reactor,首先要分清这两个核心概念的区别:

  • Mono(0 或 1 个元素的异步序列):Mono 代表一个最多只包含单个元素的异步计算结果。你可以把它理解为异步版的 Optional 或 CompletableFuture。
    典型场景:根据 ID 查询单个用户信息、保存一条记录、执行一次无返回值的异步操作(如 Mono)、HTTP 接口返回单个对象等。
  • Flux(0 到 N 个元素的异步序列):
    Flux 代表一个包含 0 到多个元素的有序异步序列,它甚至可以是一个无限流。你可以把它想象成一条物流传送带,或者数据库的游标。
    典型场景:查询用户列表、处理文件中的多行数据、WebSocket 消息流、实时传感器数据推送等。

声明式与惰性执行(Lazy Evaluation)

这是响应式编程中最反直觉、但也最核心的特性。在 Reactor 中,当你调用 map、filter 等操作符时,实际上并没有任何数据被处理,也没有任何业务逻辑被执行。这些操作仅仅是在构建一条“处理流水线(Pipeline)”。

只有当有 Subscriber(订阅者)调用 subscribe() 方法时,整条流水线才会被激活,数据才会像水流一样从源头开始向下流动。这种惰性执行机制使得我们可以像搭积木一样灵活地组装和复用数据流逻辑,同时也避免了不必要的资源消耗。

数据流的生命周期与“弹珠图(Marble Diagrams)”

doOnXXX是响应式流里的副作用(side-effect)观察者:它监听信号经过、不修改流、不改变元素,返回的是同一个流。

  • 用途:打日志、埋点、调试、资源清理。
  • 铁律:改造用map/filter,观察才用doOnXXX;不subscribe不触发(冷流)。

速查表(收藏级)

方法触发信号典型用途
doFirst订阅前(仅 1 次)一次性前置准备
doOnSubscribe订阅初始化、拿 Subscription
doOnRequest下游请求观察背压
doOnCancel取消清理
doOnNext每个元素日志 / 埋点
doOnEach所有 Signal全信号观察(调试)
doOnComplete正常结束收尾
doOnError错误错误日志(降级前可见)
doOnTerminate终止前终止前打扫
doAfterTerminate终止后终止后打扫
doOnSuccessMono 成功Mono 收尾(值可能为 null)
doFinally任意终止兜底清理(带SignalType
doOnDiscard元素被丢弃释放被丢元素持有的资源

按信号阶段分类(13 个方法)

阶段方法
订阅 / 生命周期doFirstdoOnSubscribedoOnRequestdoOnCancel
元素级doOnNextdoOnEach
终止doOnCompletedoOnErrordoOnTerminatedoAfterTerminatedoOnSuccess(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_COMPLETE
  • doOnRequest9223372036854775807=Long.MAX_VALUE,即默认subscribe无限请求(背压全开)。
  • 顺序口诀:先订阅 → 后请求 → 逐个 next/消费 → 完成前 terminate → 下发后 afterTerminate → 最后 finally
Case B · 错误路径全家桶(一次演示 5 个方法)
Flux.just(2,0).map(i->10/i)// i=0 时抛 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 >> 消费 -1
  • doOnErroronErrorResume降级之前就看到了原始异常。
  • doFinallyON_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] 空完成,v="+v))// v == null.subscribe();

输出(非空 Mono):

[doOnSubscribe] [doOnNext] hello [doOnSuccess] 成功,值=hello [doFinally] 类型=ON_COMPLETE

输出(空 Mono):

[doOnSuccess] 空完成,v=null
  • doOnSuccess=doOnNext+doOnComplete的合体(Mono 专用)。
  • 空完成时v == null,判空要小心。
Case D · 背压与丢弃(doOnRequest+doOnDiscard
Flux.range(1,100).doOnRequest(n->System.out.println("[doOnRequest] "+n)).doOnDiscard(Integer.class,i->System.out.println("[doOnDiscard] 丢弃 "+i)).onBackpressureDrop()// 下游不取就丢弃.subscribe(newBaseSubscriber<Integer>(){@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 · 取消(doOnCancel
Disposabled=Flux.interval(Duration.ofMillis(100)).doOnCancel(()->System.out.println("[doOnCancel] 被取消了")).subscribe(v->System.out.println(">> "+v));Thread.sleep(350);d.dispose();// 主动取消 → 触发 doOnCancel
Case 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在链上的位置决定它在哪条线程执行——排查"日志重复/顺序乱"时这是第一怀疑点。

坑 2:doOnXXX内抛异常会污染整条流
Flux.just(1,2).doOnNext(i->{if(i==2)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()) 直接抛出异常信号。
    // 创建空流(订阅后直接触发 onComplete)Mono.empty().subscribe(data->{},error->{},()->System.out.println("空流已完成"));// 创建错误流(订阅后直接触发 onError)Flux.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 的核心区别,并学习如何使用操作符对数据流进行转换、过滤与异常处理。