【WebFlux】第三篇 —— 常用操作符与数据流转换

📅 2026/7/20 19:44:10
【WebFlux】第三篇 —— 常用操作符与数据流转换
操作符响应式流的“加工车间”在前面的笔记中我们学会了如何创建 Mono 和 Flux 数据流。但在真实的业务场景中数据往往需要经过清洗、转换、过滤或聚合。Reactor 提供了数百个操作符Operators来完成这些工作。你可以把操作符想象成工厂流水线上的各种加工机器它们接收上游的数据流经过处理后向下游输出新的数据流。本篇笔记将带你掌握最核心、最常用的几类操作符让你能够熟练地对数据流进行加工。转换操作符map 与 flatMap 的核心对决这是响应式编程中最重要也是最容易混淆的一对操作符。map同步 1对1 转换map 用于对流中的每个元素执行一个同步的转换函数。它的输入是一个元素输出也是一个元素数据流的数量保持不变。代码示例将字符串列表转换为它们的长度。Flux.just(Java,Spring,WebFlux).map(String::length)// 同步转换字符串 - 长度.subscribe(length-System.out.println(长度: length));// 输出: 4, 6, 7flatMap异步 1对N 转换flatMap 用于将一个元素转换为一个新的异步数据流Publisher然后将这些子流扁平化Flatten合并成一个新流。它的输入是一个元素输出是一个流。代码示例将每个单词拆分成单个字符。Flux.just(Hi,OK).flatMap(word-Flux.fromArray(word.split()))// 1个单词 - 包含多个字符的子流.subscribe(System.out::println);// 输出: H, I, O, K (注意flatMap不保证子流间的绝对顺序因为它是并发处理的)concatMap保证顺序的异步转换如果你既需要 flatMap 的异步转换能力又需要保证严格的顺序可以使用 concatMap。它会等待前一个子流完全处理完毕后才开始处理下一个子流。过滤操作符精准挑选数据当数据流中混杂了不需要的元素时过滤操作符就派上用场了。filter条件过滤只保留满足条件的元素。Flux.range(1,10).filter(num-num%20)// 过滤出偶数.subscribe(even-System.out.println(偶数: even));// 输出: 2, 4, 6, 8, 10distinct去重自动过滤掉流中重复的元素。Flux.just(1,2,2,3,3,4).distinct().subscribe(System.out::println);// 输出: 1, 2, 3, 4take 与 skip截取与跳过take(n) 只获取前 n 个元素拿完即关水龙头skip(n) 则跳过前 n 个元素。Flux.just(A,B,C,D,E).skip(2)// 跳过前2个.take(2)// 只取接下来的2个.subscribe(System.out::println);// 输出: C, D组合操作符拼接多根水管在实际开发中我们经常需要把多个数据流合并成一个。Reactor 提供了两种截然不同的组合方式merge / mergeWith交错合并将多个流合并为一个谁先产生数据谁就先往下游流。FluxStringflux1Flux.just(A1,A2).delayElements(Duration.ofMillis(100));FluxStringflux2Flux.just(B1,B2);Flux.merge(flux1,flux2).subscribe(System.out::println);// 输出可能是: B1, B2, A1, A2 (因为flux2没有延迟会先发射)concat / concatWith首尾相接严格按照顺序等第一个流完全结束后才开始流第二个流的数据。Flux.concat(flux1,flux2).subscribe(System.out::println);// 输出严格为: A1, A2, B1, B2zip / zipWith配对组合像“相亲”一样把两个流中相同位置的元素配对组合成一个元组Tuple。Flux.zip(Flux.just(张三,李四),Flux.just(18,25)).subscribe(tuple-System.out.println(tuple.getT1() 年龄: tuple.getT2()));// 输出: 张三 年龄: 18, 李四 年龄: 25聚合操作符将流收敛为单值有时候我们需要对整个流进行计算最终只输出一个结果。reduce聚合计算Flux.range(1,5).reduce((a,b)-ab)// 计算 12345.subscribe(sum-System.out.println(总和: sum));// 输出: 15collectList收集为集合将流中的所有元素收集到一个 List 中返回一个 MonoList。这在需要将响应式结果传递给不支持响应式的传统代码时非常有用。Flux.just(Apple,Banana).collectList().subscribe(list-System.out.println(收集结果: list));// 输出: [Apple, Banana]本篇小结掌握 map、flatMap、filter 和 merge/concat 是响应式编程的必修课。记住一个核心原则如果转换结果是一个普通对象用 map如果转换结果是一个新的 Mono/Flux用 flatMap。下一步预告数据流在流转过程中难免会遇到异常如网络超时、数据库报错。在响应式编程中传统的 try-catch 已经失效。下一篇笔记我们将深入探讨响应式流的异常处理机制并学习如何编写优雅的响应式单元测试。