【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 的核心区别并学习如何使用操作符对数据流进行转换、过滤与异常处理。

相关新闻

商业视频监控系统智能化升级与EasyGBS平台实践

商业视频监控系统智能化升级与EasyGBS平台实践

1. 商业场所视频监控的现状与挑战现代商业场所的视频监控系统早已超越了简单的安全防范功能,正逐步演变为支撑商业运营决策的核心基础设施。作为一名在安防行业深耕多年的从业者,我见证了商业监控系统从模拟到数字、从孤立到联网、从被动录像到主动分析的…

2026/7/27 7:55:24阅读更多 →
小熊猫Dev-C++:你的第一个C++开发环境终极指南

小熊猫Dev-C++:你的第一个C++开发环境终极指南

小熊猫Dev-C:你的第一个C开发环境终极指南 【免费下载链接】Dev-CPP A greatly improved Dev-Cpp 项目地址: https://gitcode.com/gh_mirrors/dev/Dev-CPP 你是否正在寻找一款轻量级C开发环境?厌倦了复杂配置和臃肿的IDE?小熊猫Dev-C&…

2026/7/27 7:55:24阅读更多 →
AI驱动营销转型:从碎片化内容到智能品牌建设

AI驱动营销转型:从碎片化内容到智能品牌建设

1. 营销范式革命:从"大广告"到"碎片化内容"的必然转型上世纪90年代,秦池酒业以6666万元天价夺得央视标王,次年销售额直接从1亿飙升至9.5亿。这种"春晚效应"印证了中心化媒体时代的黄金法则——品牌只需集中火力…

2026/7/27 7:55:24阅读更多 →
C/C++ for循环进阶:从内存池到图像卷积的实战面试技巧

C/C++ for循环进阶:从内存池到图像卷积的实战面试技巧

1. 项目概述:为什么“玩转for循环”是2024年面试的硬通货? 如果你最近在准备C/C的面试,尤其是那些要求手撕代码或者考察项目经验的岗位,可能会发现一个有趣的现象:面试官对“for循环”的考察,早已不是问你“…

2026/7/27 9:30:22阅读更多 →
DM6435外设接口时序与寄存器配置实战指南

DM6435外设接口时序与寄存器配置实战指南

1. 项目概述与核心价值在嵌入式DSP系统开发中,尤其是面对像德州仪器TMS320DM6435这类高性能数字媒体处理器时,最让工程师头疼的往往不是算法实现,而是如何让芯片与外部世界“对话”得稳定可靠。这个“对话”的规则手册,就是各个外…

2026/7/27 9:30:22阅读更多 →
探索高效AI助手:全面掌握Goose桌面应用可视化操作指南

探索高效AI助手:全面掌握Goose桌面应用可视化操作指南

探索高效AI助手:全面掌握Goose桌面应用可视化操作指南 【免费下载链接】goose an open source, extensible AI agent that goes beyond code suggestions - install, execute, edit, and test with any LLM 项目地址: https://gitcode.com/GitHub_Trending/goose3…

2026/7/27 9:30:22阅读更多 →
深入解析SM320F28335-HT:高温工业级DSC的架构、外设与电机控制实战

深入解析SM320F28335-HT:高温工业级DSC的架构、外设与电机控制实战

1. 项目概述:为什么选择SM320F28335-HT这颗“工业硬汉”?在电机驱动、数字电源或者高精度伺服控制这类项目里摸爬滚打过的工程师,大概都经历过这样的纠结:用传统的MCU吧,面对复杂的数学运算(比如Clark/Park…

2026/7/27 9:30:22阅读更多 →
SM320F28335-HT DSP外设时序深度解析:从理论到工程实践

SM320F28335-HT DSP外设时序深度解析:从理论到工程实践

1. 项目概述与核心价值在电机驱动、数字电源或者任何需要精密时序控制的嵌入式系统里,选对了高性能的DSP只是第一步,比如德州仪器的SM320F28335-HT。这颗芯片真正的威力,藏在它的数据手册里那些密密麻麻的时序参数表格和波形图里。我见过不少…

2026/7/27 9:30:22阅读更多 →
AI论文写作工具对比:千笔与灵感AI的专科生应用指南

AI论文写作工具对比:千笔与灵感AI的专科生应用指南

1. 论文写作工具现状与需求分析作为一名在学术写作领域摸爬滚打多年的从业者,我深刻理解专科生在论文写作过程中面临的困境。时间紧、任务重、学术基础薄弱是普遍存在的三大难题。传统写作方式下,学生需要花费大量时间在文献检索、框架搭建和内容组织上&…

2026/7/27 9:28:22阅读更多 →
覆盖国产 + 海外 + 开源模型,OpenClaw 2.7.9 Windows/Mac 双端部署详解

覆盖国产 + 海外 + 开源模型,OpenClaw 2.7.9 Windows/Mac 双端部署详解

🔹 工具基础介绍 OpenClaw 是开源生态中一款实用性较强的本地智能工具,凭借本地离线运行、可视化图形操作和任务自动化三大核心特性,赢得了众多用户的青睐。与普通在线对话AI工具不同,它属于能够直接操控本机软硬件的智能数字员工…

2026/7/27 1:14:34阅读更多 →
伺服阀焊完微漏毁整机?精密激光焊接三关锁住高压

伺服阀焊完微漏毁整机?精密激光焊接三关锁住高压

所谓液压伺服阀体的精密激光焊接,是用激光束对阀座壳体(通常为不锈钢或铝合金)进行密封焊接,使阀体在21-35MPa的高压液压油或压缩气体中长期运行而不发生介质泄漏。液压伺服阀是高端液压系统的"大脑"。从航空航天飞行控…

2026/7/27 1:14:52阅读更多 →
D2DX:三步实现《暗黑破坏神2》高清宽屏体验的终极指南

D2DX:三步实现《暗黑破坏神2》高清宽屏体验的终极指南

D2DX:三步实现《暗黑破坏神2》高清宽屏体验的终极指南 【免费下载链接】d2dx D2DX is a complete solution to make Diablo II run well on modern PCs, with high fps and better resolutions. 项目地址: https://gitcode.com/gh_mirrors/d2/d2dx 你是否还在…

2026/7/27 1:14:56阅读更多 →
SPI实战指南:从时钟模式到寄存器配置,解决嵌入式通信难题

SPI实战指南:从时钟模式到寄存器配置,解决嵌入式通信难题

1. 项目概述:从寄存器手册到实战指南 如果你手头有一份类似德州仪器(TI)TMS320x240xA系列DSP的SPI模块技术手册,看着里面密密麻麻的寄存器位定义、时序图和公式,是不是感觉头大?这份资料虽然权威&#xff0…

2026/7/27 0:00:24阅读更多 →
【JAVA毕设源码分享】基于springboot的水果购物管理系统的设计与实现(程序+文档+代码讲解+一条龙定制)

【JAVA毕设源码分享】基于springboot的水果购物管理系统的设计与实现(程序+文档+代码讲解+一条龙定制)

博主介绍:✌️码农一枚 ,专注于大学生项目实战开发、讲解和毕业🚢文撰写修改等。全栈领域优质创作者,博客之星、掘金/华为云/阿里云/InfoQ等平台优质作者、专注于Java、小程序技术领域和毕业项目实战 ✌️技术范围:&am…

2026/7/27 0:00:24阅读更多 →
2007-2023年各市区县生态文明建设示范区DID

2007-2023年各市区县生态文明建设示范区DID

数据简介 自改革开放以来,我国依赖高投入、高资源消耗和高污染等传统发展模式实现了经济短期内的快速增长, 然而这也导致了严重的生态环境危机。因此,国家有力于推动企业高质量经济发展,协同生态保护的方针,从而从201…

2026/7/27 0:00:24阅读更多 →
YOLOv8推理性能优化:从1.2FPS到35FPS的全链路加速实践

YOLOv8推理性能优化:从1.2FPS到35FPS的全链路加速实践

如果你在部署 YOLOv8 时,发现推理速度只有可怜的 1-2 FPS,而别人的演示视频却能跑到 30 FPS 以上,那么问题很可能不在模型本身,而在于你的整个处理链路。很多开发者拿到一个训练好的 YOLOv8 模型后,会直接使用官方示例…

2026/7/25 23:03:25阅读更多 →
Coze与Dify对比指南:低代码AI应用开发从入门到实战

Coze与Dify对比指南:低代码AI应用开发从入门到实战

1. 从零到一:为什么你需要了解 Coze 和 Dify?如果你对 AI 应用开发感兴趣,但一看到“大模型”、“智能体”、“工作流”这些词就头疼,觉得门槛太高,那这篇文章就是为你准备的。很多开发者,包括我自己&#…

2026/7/26 19:05:21阅读更多 →
AI生图工具怎么选?2026年6月版实测对比

AI生图工具怎么选?2026年6月版实测对比

做自媒体的朋友应该都有体会:配图一直是个让人头疼的问题。2026年,AI生图工具已经非常成熟了,但工具太多反而不知道怎么选。以下是截至2026年6月我对主流AI生图工具的实测对比。Midjourney V8.1:速度之王2026年6月11日&#xff0c…

2026/7/26 19:05:21阅读更多 →