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查询单个用户实体:
Mono<User> findById(String id) - 执行一个更新操作,返回更新成功的记录数:
Mono<Long> update(User user) - 发起一个HTTP请求并等待单个响应体:
Mono<HttpResponse> - 任何返回
Optional<T>的地方,都可以自然地转换为Mono<T>(甚至可以是空的Mono.empty())。
关键理解:即使操作本身是异步非阻塞的,但Mono所承载的“结果”的预期是单一的。它是对一次异步计算结果的封装。
2.2 Flux:处理数据流
Flux则是真正的“流”的化身。它代表一个异步的、可以发射0到N个数据元素的序列, optionally followed by a completion signal or an error.
核心语义:0到N个元素。典型场景:
- 查询所有用户列表:
Flux<User> findAll() - 监听一个消息队列(如Kafka topic),持续消费消息:
Flux<Message> streamFromKafka() - 服务器发送事件(Server-Sent Events, SSE):客户端保持连接,服务器持续推送事件流。
- 读取一个大文件,按行处理:
Flux<String> lines = Flux.using(...)
关键理解:Flux处理的是可能无界(无限)的数据流。你需要使用操作符(Operators)来声明对这个流中每个元素的处理逻辑,而不是一次性获取所有数据。
2.3 二者关系与选择
- 包含关系:从数据量的角度看,
Mono可以被视为Flux的一种特例(最多一个元素)。事实上,很多操作符是通用的,且Mono可以轻松转换为Flux(例如mono.flux())。但反之,将Flux转为期待单个结果的Mono则需要聚合操作,如collectList()、reduce()。 - 选择依据:选择的关键在于业务语义,而非性能。如果你明确知道操作的结果是单个对象(或可能没有),就用
Mono。如果你要处理一个集合或一个持续的数据流,就用Flux。错误的选用(比如用Flux返回单个对象)会让API的使用者感到困惑,不知道你到底想返回什么。
实操心得:在定义Repository或Service层接口时,我强烈建议根据方法语义严格区分返回类型。
findById返回Mono<User>,findAll返回Flux<User>。这本身就是一种清晰的设计文档。刚开始可能会觉得多此一举,但这对团队协作和代码可读性有巨大好处。
3. 创建数据流:多种源头,声明式起点
响应式编程的第一步是创建数据源。Reactor提供了极其丰富和灵活的静态工厂方法来创建Flux和Mono。
3.1 简单静态创建
这是最直接的创建方式,适用于已知的、有限的数据。
// 1. 从已知值创建 Mono<String> mono = Mono.just("Hello"); Flux<String> flux = Flux.just("A", "B", "C"); // 2. 从Optional或可能为null的值创建 (避免NPE的神器) Mono<String> monoFromNullable = Mono.justOrEmpty(somePotentiallyNullString); // 如果somePotentiallyNullString是null,则创建一个空的Mono(不发射元素,直接完成)。 // 3. 创建空流或错误流 Mono<Void> emptyMono = Mono.empty(); Flux<String> emptyFlux = Flux.empty(); Mono<Object> errorMono = Mono.error(new RuntimeException("Oops!")); Flux<Object> errorFlux = Flux.error(new RuntimeException("Stream failed!")); // 4. 从Iterable、数组、Stream创建 List<String> list = Arrays.asList("x", "y", "z"); Flux<String> fluxFromIterable = Flux.fromIterable(list); Flux<String> fluxFromArray = Flux.fromArray(new String[]{"a", "b"}); Flux<String> fluxFromStream = Flux.fromStream(list.stream()); // 重要:来自Stream时需注意关闭问题,建议使用Flux.using或确保流被正确管理。3.2 动态与区间创建
当数据有规律或需要动态生成时,这些方法非常有用。
// 1. 生成数字范围 Flux<Integer> range = Flux.range(1, 5); // 发射 1,2,3,4,5 // 2. 间隔生成,用于模拟心跳、定时任务 Flux<Long> interval = Flux.interval(Duration.ofSeconds(1)); // 从0开始,每秒发射一个递增的Long。这是一个无限流! // 测试时一定要用take等操作符限制,否则程序不会结束。 Flux<Long> fiveTicks = Flux.interval(Duration.ofMillis(500)).take(5); // 只取前5个 // 3. 使用generate创建有状态的流 (类似Stream.generate) // 同步、逐一生成,可以基于上一个状态计算下一个 Flux<Integer> statefulFlux = Flux.generate( () -> 0, // 初始状态 (state, sink) -> { sink.next(state); // 发射当前状态 if (state == 10) { sink.complete(); // 完成流 } return state + 1; // 返回新状态 } );3.3 高级创建:create与push
create和push是更强大的方法,允许你将一个现有的、基于回调或事件的API(通常是多线程的)桥接到响应式世界。这是集成非响应式库的关键。
// 使用create,它可以处理多线程发射 Flux<String> 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 过滤与筛选
用于从流中挑选出需要的元素。
Flux<Integer> numbers = Flux.range(1, 10); // filter: 过滤 Flux<Integer> even = numbers.filter(i -> i % 2 == 0); // 2,4,6,8,10 // distinct: 去重 Flux<String> withDuplicates = Flux.just("a", "b", "a", "c"); Flux<String> distinct = withDuplicates.distinct(); // a, b, c // take: 取前N个 / takeLast / takeWhile / takeUntil Flux<Integer> firstThree = numbers.take(3); // 1,2,3 Flux<Integer> whileLessThanFive = numbers.takeWhile(i -> i < 5); // 1,2,3,4 // skip: 跳过前N个 Flux<Integer> skipFirstThree = numbers.skip(3); // 4,5,6,7,8,9,104.2 映射与转换
将流中的每个元素转换为另一种形式。
Flux<String> names = Flux.just("alice", "bob"); // map: 一对一同步转换 Flux<String> upper = names.map(String::toUpperCase); // ALICE, BOB Flux<Integer> lengths = names.map(String::length); // 5, 3 // flatMap: 一对多异步转换 (最强大也最常用) // 将每个元素映射为一个新的Publisher(Flux/Mono),然后将所有这些Publisher“扁平化”成一个新的Flux。 Flux<User> users = Flux.just("user1", "user2"); Flux<Order> orders = users.flatMap(userId -> orderRepository.findOrdersByUserId(userId) // 假设返回 Flux<Order> ); // 如果orderRepository.findOrdersByUserId(userId)对user1返回Order[A,B],对user2返回Order[C] // 那么最终的orders流发射顺序可能是:A, B, C(注意:由于异步,顺序可能交错)。 // concatMap: 类似flatMap,但会保证内部Publisher的顺序,一个接一个处理,不交错。性能不如flatMap,但有序。 // switchMap: 更特殊,当源发射新元素时,会取消上一个元素映射出的Publisher的订阅。常用于搜索联想词等场景。4.3 组合与聚合
将多个流组合,或将一个流聚合成一个结果(通常是Mono)。
Flux<Integer> flux1 = Flux.just(1, 2, 3); Flux<Integer> flux2 = Flux.just(4, 5, 6); // merge: 合并多个流,元素按实际到达时间交错发射 Flux<Integer> merged = Flux.merge(flux1, flux2); // 可能顺序:1,4,2,5,3,6 // concat: 连接多个流,先发射第一个流的全部元素,再发射第二个流的全部元素 Flux<Integer> concatenated = Flux.concat(flux1, flux2); // 保证顺序:1,2,3,4,5,6 // zip: 将多个流中的元素一对一配对组合(像拉链一样) Flux<String> names = Flux.just("Alice", "Bob"); Flux<Integer> ages = Flux.just(30, 25); Flux<Tuple2<String, Integer>> zipped = Flux.zip(names, ages); // 发射:("Alice", 30), ("Bob", 25) // 如果流长度不同,以最短的流结束为准。 // reduce 和 scan: 聚合 Mono<Integer> sumMono = Flux.range(1, 5).reduce(0, (a, b) -> a + b); // Mono发射 15 Flux<Integer> scanFlux = Flux.range(1, 5).scan(0, (a, b) -> a + b); // Flux发射 0,1,3,6,10,15 // reduce返回最终结果的Mono,scan返回每个中间结果的Flux。 // collectList / collectSortedList: 将Flux中的所有元素收集到一个List中,返回Mono<List<T>> Mono<List<Integer>> listMono = Flux.range(1, 3).collectList();4.4 错误处理
响应式流中的错误是一个终止信号。一旦发生错误,流就会停止。我们必须妥善处理。
Flux<Integer> flux = Flux.range(1, 5) .map(i -> { if (i == 3) throw new RuntimeException("Boom on 3!"); return i * 2; }); // 1. onErrorReturn: 发生错误时,返回一个静态默认值 Flux<Integer> withDefault = flux.onErrorReturn(0); // 发射:2, 4, 0 (遇到错误,返回0,然后流完成) // 2. onErrorResume: 发生错误时,切换到一个备用的Publisher Flux<Integer> withFallback = flux.onErrorResume(e -> Flux.just(9, 10, 11)); // 发射:2, 4, 9, 10, 11 // 3. onErrorContinue: 一个危险但有时有用的操作符。它允许错误发生后,丢弃出错的元素,但继续处理流中的后续元素。 // 使用需极度小心,因为它改变了正常的错误传播语义。 Flux<Integer> 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() { Flux<String> flux = Flux.just("foo", "bar"); StepVerifier.create(flux) .expectNext("foo") // 期待下一个元素是"foo" .expectNext("bar") // 然后是"bar" .verifyComplete(); // 最后期待流正常完成 } @Test void testMonoError() { Mono<Object> 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 Mono<User> getUser(@PathVariable String id) { return userRepository.findById(id); // Spring会自动处理订阅和响应序列化 } // 返回资源集合流 @GetMapping public Flux<User> getAllUsers() { return userRepository.findAll(); // 数据会以流的形式(如NDJSON)逐步发送到客户端 } // 处理Server-Sent Events @GetMapping(value = "/events", produces = MediaType.TEXT_EVENT_STREAM_VALUE) public Flux<UserEvent> getUserEvents() { return userRepository.findUserEvents(); // 返回一个持续的事件流 } // 接收请求体流并处理 @PostMapping public Mono<User> createUser(@RequestBody Mono<User> 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 Mono<User> 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 Flux<Post> 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 ReactiveCrudRepository<User, String> { @Query("SELECT * FROM users WHERE age > :age") Flux<User> findByAgeGreaterThan(int age); Mono<User> 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返回Mono<T>的方法内部),调用了block()等待结果,这完全违背了响应式的初衷,会阻塞负责处理请求的线程(如Netty的EventLoop),导致性能灾难。
// 错误做法 public Mono<User> getUserWrong(String id) { User user = reactiveRepo.findById(id).block(); // 阻塞! return Mono.just(user); } // 正确做法 public Mono<User> getUserRight(String id) { return reactiveRepo.findById(id); // 直接返回Mono }陷阱三:flatMap的并发控制flatMap默认的并发度是256(Queues.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是你的好朋友),多观察日志,谨慎处理阻塞和错误,响应式编程的道路就会越走越顺。