1. 项目概述从“拉”到“推”的思维跃迁如果你是从传统的Spring MVC或者Servlet体系转过来的开发者第一次接触“响应式”这个词大概率会有点懵。我们习惯了“拉”取数据——发一个HTTP请求线程阻塞等待数据库返回结果然后组装、渲染、返回。整个链路清晰但效率瓶颈也显而易见每个阻塞的线程都在占用宝贵的系统资源尤其是在高并发、高延迟的I/O场景下线程池被打满是分分钟的事。“05-流式操作使用 Flux 和 Mono 构建响应式数据流”这个标题直指响应式编程的核心武器库。它不是一个简单的API教程而是一套全新的异步非阻塞编程范式。Flux和 Mono作为Project Reactor库也是Spring WebFlux的基石中的两个核心发布者Publisher代表的就是这种“推”的模式。数据不再是你去数据库里“捞”出来的而是作为一个潜在的、可能无限的数据流由源头“推送”给你。你的任务是如何优雅地、高效地处理这个流无论是处理其中零个、一个还是成千上万个元素。这背后的驱动力是现实应用对高吞吐量和低资源消耗的极致追求。想象一个实时股票报价系统、一个物联网平台处理海量传感器数据或者一个社交媒体的消息推送流。传统阻塞模型在这里会迅速崩溃而基于Flux和Mono的响应式流则能像铺设好的高速公路一样让数据元素作为“车辆”持续、异步地通过系统资源尤其是线程得到最大化复用。理解并掌握它们不仅是学习几个新类更是将你的后端开发思维从“同步世界”升级到“异步世界”的关键一步。2. 核心概念解析Flux与Mono的二元世界在深入操作之前我们必须先厘清Flux和Mono的根本区别与联系。这是所有响应式操作的起点理解错了后面的代码就会写得别别扭扭。2.1 Mono期待单次结果你可以把Mono理解为Future或CompletableFuture在响应式世界的对应物但它更强大、更声明式。它代表一个异步的、最多发射一个数据元素或一个错误的序列。核心语义0或1个元素。典型场景根据ID查询单个用户实体MonoUser findById(String id)执行一个更新操作返回更新成功的记录数MonoLong update(User user)发起一个HTTP请求并等待单个响应体MonoHttpResponse任何返回OptionalT的地方都可以自然地转换为MonoT甚至可以是空的Mono.empty()。关键理解即使操作本身是异步非阻塞的但Mono所承载的“结果”的预期是单一的。它是对一次异步计算结果的封装。2.2 Flux处理数据流Flux则是真正的“流”的化身。它代表一个异步的、可以发射0到N个数据元素的序列 optionally followed by a completion signal or an error.核心语义0到N个元素。典型场景查询所有用户列表FluxUser findAll()监听一个消息队列如Kafka topic持续消费消息FluxMessage streamFromKafka()服务器发送事件Server-Sent Events, SSE客户端保持连接服务器持续推送事件流。读取一个大文件按行处理FluxString lines Flux.using(...)关键理解Flux处理的是可能无界无限的数据流。你需要使用操作符Operators来声明对这个流中每个元素的处理逻辑而不是一次性获取所有数据。2.3 二者关系与选择包含关系从数据量的角度看Mono可以被视为Flux的一种特例最多一个元素。事实上很多操作符是通用的且Mono可以轻松转换为Flux例如mono.flux()。但反之将Flux转为期待单个结果的Mono则需要聚合操作如collectList()、reduce()。选择依据选择的关键在于业务语义而非性能。如果你明确知道操作的结果是单个对象或可能没有就用Mono。如果你要处理一个集合或一个持续的数据流就用Flux。错误的选用比如用Flux返回单个对象会让API的使用者感到困惑不知道你到底想返回什么。实操心得在定义Repository或Service层接口时我强烈建议根据方法语义严格区分返回类型。findById返回MonoUserfindAll返回FluxUser。这本身就是一种清晰的设计文档。刚开始可能会觉得多此一举但这对团队协作和代码可读性有巨大好处。3. 创建数据流多种源头声明式起点响应式编程的第一步是创建数据源。Reactor提供了极其丰富和灵活的静态工厂方法来创建Flux和Mono。3.1 简单静态创建这是最直接的创建方式适用于已知的、有限的数据。// 1. 从已知值创建 MonoString mono Mono.just(Hello); FluxString flux Flux.just(A, B, C); // 2. 从Optional或可能为null的值创建 (避免NPE的神器) MonoString monoFromNullable Mono.justOrEmpty(somePotentiallyNullString); // 如果somePotentiallyNullString是null则创建一个空的Mono不发射元素直接完成。 // 3. 创建空流或错误流 MonoVoid emptyMono Mono.empty(); FluxString emptyFlux Flux.empty(); MonoObject errorMono Mono.error(new RuntimeException(Oops!)); FluxObject errorFlux Flux.error(new RuntimeException(Stream failed!)); // 4. 从Iterable、数组、Stream创建 ListString list Arrays.asList(x, y, z); FluxString fluxFromIterable Flux.fromIterable(list); FluxString fluxFromArray Flux.fromArray(new String[]{a, b}); FluxString fluxFromStream Flux.fromStream(list.stream()); // 重要来自Stream时需注意关闭问题建议使用Flux.using或确保流被正确管理。3.2 动态与区间创建当数据有规律或需要动态生成时这些方法非常有用。// 1. 生成数字范围 FluxInteger range Flux.range(1, 5); // 发射 1,2,3,4,5 // 2. 间隔生成用于模拟心跳、定时任务 FluxLong interval Flux.interval(Duration.ofSeconds(1)); // 从0开始每秒发射一个递增的Long。这是一个无限流 // 测试时一定要用take等操作符限制否则程序不会结束。 FluxLong fiveTicks Flux.interval(Duration.ofMillis(500)).take(5); // 只取前5个 // 3. 使用generate创建有状态的流 (类似Stream.generate) // 同步、逐一生成可以基于上一个状态计算下一个 FluxInteger statefulFlux Flux.generate( () - 0, // 初始状态 (state, sink) - { sink.next(state); // 发射当前状态 if (state 10) { sink.complete(); // 完成流 } return state 1; // 返回新状态 } );3.3 高级创建create与pushcreate和push是更强大的方法允许你将一个现有的、基于回调或事件的API通常是多线程的桥接到响应式世界。这是集成非响应式库的关键。// 使用create它可以处理多线程发射 FluxString bridge Flux.create(sink - { // 假设我们有一个传统的、基于监听器的消息源 MyMessageSource source new MyMessageSource(); source.registerListener(new MessageListener() { Override public void onMessage(String message) { sink.next(message); // 在监听器回调中发射数据 } Override public void onError(Throwable t) { sink.error(t); } Override public void onComplete() { sink.complete(); } }); // sink.onCancel 可以用于取消注册监听器进行资源清理 sink.onCancel(() - source.shutdown()); }); // push 是 create 的变体适用于单线程生产者场景性能稍好。注意事项create和push是“桥接”方法也是容易出错的地方。你必须仔细考虑背压Backpressure问题。默认情况下create使用的FluxSink是IGNORE背压策略的如果生产者太快可能导致内存溢出。对于快速生产者应考虑使用BUFFER、DROP或LATEST策略或者更好的方式是让你的消息源本身支持响应式流协议。4. 流式操作符声明式数据处理的基石操作符是响应式编程的灵魂。它们以声明式的方式串联起来形成一个处理流水线。理解操作符的分类和用途至关重要。4.1 过滤与筛选用于从流中挑选出需要的元素。FluxInteger numbers Flux.range(1, 10); // filter: 过滤 FluxInteger even numbers.filter(i - i % 2 0); // 2,4,6,8,10 // distinct: 去重 FluxString withDuplicates Flux.just(a, b, a, c); FluxString distinct withDuplicates.distinct(); // a, b, c // take: 取前N个 / takeLast / takeWhile / takeUntil FluxInteger firstThree numbers.take(3); // 1,2,3 FluxInteger whileLessThanFive numbers.takeWhile(i - i 5); // 1,2,3,4 // skip: 跳过前N个 FluxInteger skipFirstThree numbers.skip(3); // 4,5,6,7,8,9,104.2 映射与转换将流中的每个元素转换为另一种形式。FluxString names Flux.just(alice, bob); // map: 一对一同步转换 FluxString upper names.map(String::toUpperCase); // ALICE, BOB FluxInteger lengths names.map(String::length); // 5, 3 // flatMap: 一对多异步转换 (最强大也最常用) // 将每个元素映射为一个新的PublisherFlux/Mono然后将所有这些Publisher“扁平化”成一个新的Flux。 FluxUser users Flux.just(user1, user2); FluxOrder orders users.flatMap(userId - orderRepository.findOrdersByUserId(userId) // 假设返回 FluxOrder ); // 如果orderRepository.findOrdersByUserId(userId)对user1返回Order[A,B]对user2返回Order[C] // 那么最终的orders流发射顺序可能是A, B, C注意由于异步顺序可能交错。 // concatMap: 类似flatMap但会保证内部Publisher的顺序一个接一个处理不交错。性能不如flatMap但有序。 // switchMap: 更特殊当源发射新元素时会取消上一个元素映射出的Publisher的订阅。常用于搜索联想词等场景。4.3 组合与聚合将多个流组合或将一个流聚合成一个结果通常是Mono。FluxInteger flux1 Flux.just(1, 2, 3); FluxInteger flux2 Flux.just(4, 5, 6); // merge: 合并多个流元素按实际到达时间交错发射 FluxInteger merged Flux.merge(flux1, flux2); // 可能顺序1,4,2,5,3,6 // concat: 连接多个流先发射第一个流的全部元素再发射第二个流的全部元素 FluxInteger concatenated Flux.concat(flux1, flux2); // 保证顺序1,2,3,4,5,6 // zip: 将多个流中的元素一对一配对组合像拉链一样 FluxString names Flux.just(Alice, Bob); FluxInteger ages Flux.just(30, 25); FluxTuple2String, Integer zipped Flux.zip(names, ages); // 发射(Alice, 30), (Bob, 25) // 如果流长度不同以最短的流结束为准。 // reduce 和 scan: 聚合 MonoInteger sumMono Flux.range(1, 5).reduce(0, (a, b) - a b); // Mono发射 15 FluxInteger scanFlux Flux.range(1, 5).scan(0, (a, b) - a b); // Flux发射 0,1,3,6,10,15 // reduce返回最终结果的Monoscan返回每个中间结果的Flux。 // collectList / collectSortedList: 将Flux中的所有元素收集到一个List中返回MonoListT MonoListInteger listMono Flux.range(1, 3).collectList();4.4 错误处理响应式流中的错误是一个终止信号。一旦发生错误流就会停止。我们必须妥善处理。FluxInteger flux Flux.range(1, 5) .map(i - { if (i 3) throw new RuntimeException(Boom on 3!); return i * 2; }); // 1. onErrorReturn: 发生错误时返回一个静态默认值 FluxInteger withDefault flux.onErrorReturn(0); // 发射2, 4, 0 (遇到错误返回0然后流完成) // 2. onErrorResume: 发生错误时切换到一个备用的Publisher FluxInteger withFallback flux.onErrorResume(e - Flux.just(9, 10, 11)); // 发射2, 4, 9, 10, 11 // 3. onErrorContinue: 一个危险但有时有用的操作符。它允许错误发生后丢弃出错的元素但继续处理流中的后续元素。 // 使用需极度小心因为它改变了正常的错误传播语义。 FluxInteger continued Flux.just(1,2,3,4,5) .map(i - { if (i 3) throw new RuntimeException(Bad 3); return i * 10; }) .onErrorContinue((e, obj) - System.err.println(Dropped obj due to e)); // 可能发射10, 20, 40, 50 3被丢弃流继续 // 4. doOnError: 副作用操作符用于记录日志、发监控等不影响流本身。 flux.doOnError(e - log.error(Error in stream, e));实操心得错误处理策略的选择取决于业务场景。对于“尽力而为”的非关键操作如记录辅助日志失败onErrorResume或onErrorReturn很合适。对于关键业务链你可能希望错误向上传播在最外层如Controller的异常处理器统一处理。尽量避免滥用onErrorContinue除非你非常清楚自己在做什么因为它会掩盖错误使调试变得困难。5. 背压处理流控的艺术背压Backpressure是响应式编程的核心概念也是区别于传统拉模式的关键。它解决的是“生产者太快消费者太慢”的问题。在响应式流规范中订阅者Subscriber可以主动向发布者Publisher请求特定数量的数据通过Subscription.request(n)从而实现流量控制。5.1 背压策略Reactor为一些不支持背压的源头如create或需要缓冲的场景提供了不同的背压策略。Flux.create(sink - { // 一个快速的生产者 for (int i 0; i 1000000; i) { sink.next(i); } sink.complete(); }, FluxSink.OverflowStrategy.BUFFER); // 指定背压策略常见的OverflowStrategy有BUFFER默认对于某些源。如果下游跟不上将元素缓冲在内存队列中。风险是可能导致OutOfMemoryError。DROP如果下游未准备好接收直接丢弃新产生的元素。LATEST只保留最新的元素下游请求时给予最后一个。ERROR当下游跟不上时直接发出IllegalStateException错误信号。IGNORE完全忽略下游的请求可能压垮下游。Flux.create的默认策略就是它所以使用create时要特别小心。5.2 操作符与背压许多操作符内部已经智能地处理了背压。例如limitRate(n)在请求上游时进行批量化处理。下游请求N个它可能向上游请求N个但在传递给下游时会分批比如每批n/3个传递平滑流量。onBackpressureBuffer()显式添加一个无界或有界缓冲区。onBackpressureDrop()当下游跟不上时丢弃元素。onBackpressureLatest()类似LATEST策略。最佳实践对于你自己创建的流尤其是桥接外部事件源时要慎重选择背压策略。对于消费来自网络或数据库的响应式流它们通常支持背压你一般不需要显式设置操作符链会自动传播背压请求。6. 线程与调度器控制执行的上下文响应式流本身是并发无关的但它允许你通过调度器Scheduler来指定操作在哪个线程上执行。6.1 常用调度器Schedulers.immediate()在当前线程执行默认。Schedulers.single()一个专用的单线程所有任务排队执行。Schedulers.elastic()已弃用 /Schedulers.boundedElastic()适用于阻塞性I/O操作的线程池。它会根据需要创建新线程但有上限适合处理会阻塞当前线程的任务如调用一个传统的阻塞式HTTP客户端。重要不要用它执行非阻塞或计算密集型任务会造成资源浪费。Schedulers.parallel()固定大小的线程池适用于CPU密集型并行处理。大小通常等于CPU核心数。Schedulers.fromExecutorService()包装现有的ExecutorService。6.2 切换执行上下文使用publishOn和subscribeOn操作符。// subscribeOn: 指定上游源source在哪个调度器上执行。影响整个链的起点。 Mono.fromCallable(() - { // 这是一个阻塞调用 return blockingHttpClient.call(); }) .subscribeOn(Schedulers.boundedElastic()) // 将阻塞调用转移到弹性线程池 .map(response - process(response)) // 后续操作默认在发出信号的线程可能是主线程执行 .block(); // publishOn: 影响其下游操作符的执行上下文。 Flux.range(1, 10) .map(i - i * 2) // 在原始线程执行 .publishOn(Schedulers.parallel()) // 切换上下文 .filter(i - i % 3 0) // 在parallel线程池执行 .map(i - i * 10) // 在parallel线程池执行 .subscribe();关键区别subscribeOn用于改变源头source的订阅过程以及后续直到第一个publishOn之前的执行线程。整个链通常只需要一个subscribeOn且位置不限但一般靠近源头。publishOn用于改变其下游操作的执行线程。链中可以有多个publishOn每次都会切换上下文。踩坑记录最大的坑就是将阻塞代码放在非弹性调度器上执行或者反过来。我曾经在Schedulers.parallel()中执行了一个JDBC查询导致整个并行线程池被卡住应用响应性急剧下降。黄金法则阻塞操作用Schedulers.boundedElastic()非阻塞/计算操作用Schedulers.parallel()或默认。7. 测试响应式流StepVerifier与调试测试响应式流与测试普通代码不同因为它是异步的、基于事件的。Reactor提供了强大的StepVerifier工具。7.1 使用StepVerifier进行断言Test void testFlux() { FluxString flux Flux.just(foo, bar); StepVerifier.create(flux) .expectNext(foo) // 期待下一个元素是foo .expectNext(bar) // 然后是bar .verifyComplete(); // 最后期待流正常完成 } Test void testMonoError() { MonoObject errorMono Mono.error(new IllegalArgumentException(bad arg)); StepVerifier.create(errorMono) .expectErrorMatches(throwable - throwable instanceof IllegalArgumentException throwable.getMessage().equals(bad arg) ) .verify(); } Test void testWithVirtualTime() { // 测试时间相关的操作如interval无需真实等待 StepVerifier.withVirtualTime(() - Flux.interval(Duration.ofHours(1)).take(2) ) .expectSubscription() .thenAwait(Duration.ofHours(2)) // 虚拟时间快进2小时 .expectNext(0L, 1L) .verifyComplete(); }7.2 调试技巧检查日志与Hooks响应式链式调用一旦出错栈跟踪可能指向内部框架代码难以定位问题源头。启用调试模式在应用启动时设置Hooks.onOperatorDebug()。这会捕获组装操作符时的堆栈信息代价是性能开销仅用于开发环境。使用log()操作符在链中插入.log()可以打印出生命周期事件onSubscribe, request, onNext, onComplete, onError。Flux.range(1, 3) .map(i - i * 2) .log(my-stream) .subscribe();使用doOnXXX系列操作符在关键节点添加副作用回调打印信息。flux.doOnSubscribe(s - log.info(Subscribed: {}, s)) .doOnNext(i - log.info(Next: {}, i)) .doOnError(e - log.error(Error: , e)) .doOnComplete(() - log.info(Completed)) .subscribe();8. 实战集成Spring WebFlux中的Flux与Mono在Spring WebFlux中Flux和Mono直接作为HTTP请求和响应的载体实现了真正的端到端非阻塞。8.1 声明响应式控制器RestController RequestMapping(/users) public class UserController { Autowired private ReactiveUserRepository userRepository; // 返回Flux/Mono的仓库 // 返回单个资源 GetMapping(/{id}) public MonoUser getUser(PathVariable String id) { return userRepository.findById(id); // Spring会自动处理订阅和响应序列化 } // 返回资源集合流 GetMapping public FluxUser getAllUsers() { return userRepository.findAll(); // 数据会以流的形式如NDJSON逐步发送到客户端 } // 处理Server-Sent Events GetMapping(value /events, produces MediaType.TEXT_EVENT_STREAM_VALUE) public FluxUserEvent getUserEvents() { return userRepository.findUserEvents(); // 返回一个持续的事件流 } // 接收请求体流并处理 PostMapping public MonoUser createUser(RequestBody MonoUser userMono) { return userMono.flatMap(userRepository::save); } }8.2 响应式WebClient消费外部API也同样是非阻塞的。Service public class ExternalService { private final WebClient webClient; public ExternalService(WebClient.Builder builder) { this.webClient builder.baseUrl(https://api.example.com).build(); } public MonoUser fetchUser(String id) { return webClient.get() .uri(/users/{id}, id) .retrieve() .bodyToMono(User.class) .timeout(Duration.ofSeconds(5)) // 超时控制 .onErrorResume(e - Mono.just(new User(fallback))); // 降级 } public FluxPost streamPosts() { return webClient.get() .uri(/posts/stream) .accept(MediaType.TEXT_EVENT_STREAM) .retrieve() .bodyToFlux(Post.class); } }8.3 数据库集成R2DBC对于关系型数据库可以使用R2DBC驱动如r2dbc-postgresql,r2dbc-mysql配合Spring Data R2DBC。public interface ReactiveUserRepository extends ReactiveCrudRepositoryUser, String { Query(SELECT * FROM users WHERE age :age) FluxUser findByAgeGreaterThan(int age); MonoUser findByName(String name); }使用时所有方法都返回Flux或Mono数据库查询操作不会阻塞操作线程。9. 性能调优与常见陷阱经过几个项目的实战我积累了一些关于性能和使用陷阱的经验。9.1 性能调优要点调度器选择重申一遍阻塞操作JDBC、同步HTTP调用、文件IO务必使用Schedulers.boundedElastic()。CPU密集型计算使用Schedulers.parallel()。错误的选择是性能问题的首要元凶。避免在响应式链中阻塞这是死罪。永远不要在map、flatMap、filter等操作符的函数中调用Thread.sleep()、block()或者任何会阻塞线程的方法。这会卡住整个线程破坏响应性。背压与缓冲对于高速生产源如果消费者确实较慢评估使用onBackpressureBuffer的容量。无界缓冲是危险的。考虑使用limitRate来平滑请求。冷发布者与热发布者冷发布者每次订阅都会重新开始数据生产。例如Flux.range(1,10)或从数据库查询返回的Flux。这是最常见的。热发布者数据生产与订阅无关订阅者只能收到订阅后产生的数据。例如Flux.share()或Flux.replay()后的流。在需要多个订阅者共享同一实时数据源时使用但要小心资源泄漏。资源清理使用Flux.using或doFinally、doOnCancel来确保资源如数据库连接、文件句柄在流终止完成、错误、取消时被正确释放。9.2 常见陷阱与解决方案陷阱一忘记订阅这是新手最常犯的错误。Flux和Mono是声明式的蓝图只有调用subscribe()、block()测试用或被框架如WebFlux订阅时才会真正开始执行。Flux.just(1,2,3).map(i - i*2); // 什么也不会发生 Flux.just(1,2,3).map(i - i*2).subscribe(); // 开始执行陷阱二在非阻塞链中误用block()在应该是非阻塞的方法中如Controller的GetMapping返回MonoT的方法内部调用了block()等待结果这完全违背了响应式的初衷会阻塞负责处理请求的线程如Netty的EventLoop导致性能灾难。// 错误做法 public MonoUser getUserWrong(String id) { User user reactiveRepo.findById(id).block(); // 阻塞 return Mono.just(user); } // 正确做法 public MonoUser getUserRight(String id) { return reactiveRepo.findById(id); // 直接返回Mono }陷阱三flatMap的并发控制flatMap默认的并发度是256Queues.SMALL_BUFFER_SIZE。这意味着对于上游的每个元素flatMap会立即订阅其内部Publisher最多同时有256个内部Publisher在运行。如果内部任务是IO密集型且数量巨大如调用十万次外部API这可能会瞬间打垮下游服务或耗尽资源。Flux.range(1, 100000) .flatMap(id - callExternalApi(id), 10) // 将并发度限制为10 .subscribe();使用flatMap的重载方法指定并发参数或者使用concatMap顺序执行并发度为1来限制并发。陷阱四错误处理吞没异常过于宽泛的onErrorResume可能会吞掉本应暴露出来的致命异常使得调试极其困难。flux.onErrorResume(e - Mono.empty()); // 任何错误都静默返回空太危险了至少应该记录日志并根据异常类型区别处理。陷阱五线程上下文丢失由于操作符切换线程ThreadLocal中的信息如Spring Security的SecurityContext、MDC日志跟踪ID可能会丢失。需要使用contextWrite操作符来传播上下文。Mono.deferContextual(ctx - { String traceId ctx.get(TRACE_ID); return makeRequestWithTraceId(traceId); }) .contextWrite(Context.of(TRACE_ID, 12345));掌握Flux和Mono的流式操作是一个从“会写”到“写好”的持续过程。它要求开发者从被动的、同步的思维模式转向主动的、声明式的、异步的思维模式。开始时可能会觉得抽象和容易出错但一旦熟悉了这种“数据流管道”的构建方式你就会发现它在处理复杂异步逻辑、构建高性能系统方面的巨大威力。记住多写测试StepVerifier是你的好朋友多观察日志谨慎处理阻塞和错误响应式编程的道路就会越走越顺。