写响应式代码的人大概都有过这种体验:Flux 和 Mono 表面上就是把数据包成一条流,可真到了要把两条流合在一起、或者想让流稍微“慢半拍”的时候,组合操作符和延迟操作符的选择立刻变得不直观了。我在一个 Spring WebFlux 网关项目里就因为在“等前一批跑完”和“谁先到谁先走”之间选错了操作符,上线第一天订单推送顺序乱掉,被迫紧急热修复。后来我把这些操作符的底层行为、调度器选择和适用场景彻底过了一遍,这篇就当作一次完整的复盘。
这篇内容主要讲两块:一是组合操作符(concatWith、mergeWith、zip 以及它们的变体),二是延迟操作符(delayElements、delaySubscription、delaySequence、delayUntil),最后落到实际聚合接口的并行化改造上。无论你是在写聚合 API、做网关转发,还是刚接触响应式想少踩几个坑,这篇都值得看完。
1. 组合操作符:concatWith、mergeWith、zip 的底层行为差异
很多人第一眼看到这三个操作符,会觉得它们做的事情都一样:把两个流合并成一个。实际上它们的差异非常大,用错之后的表现也不一样——有的是顺序错乱,有的是数据丢失,有的是整个流被最慢的一路“锁死”。我自己的经验是先分清三者的语义模型,再记用法,否则代码写对了也是一头雾水。
1.1 concatWith:接力棒式的顺序连接
concatWith的语义是:先订阅当前 Flux,等它完整结束后,再订阅传入的另一个 Publisher,然后把第二个流的数据接在第一个流后面。
Flux<Integer> first = Flux.just(1, 2, 3); Flux<Integer> second = Flux.just(4, 5, 6); first.concatWith(second) .subscribe(System.out::println); // 输出: // 1 // 2 // 3 // 4 // 5 // 6注意一个容易被忽略的细节:concatWith不会在创建时就去订阅第二个流,而是在第一个流onComplete之后才开始订阅。这意味着如果第二个流是热数据源(比如一个Sinks.Many或者已经产生了一段时间的Flux.interval),那么第一个流运行期间的这部分数据你是收不到的。
这一点很重要。很多人在做“先后两次请求”或者“先查缓存再查数据库”时觉得用 concatWith 很自然,但它本质上只是顺序连接,不是条件判断。如果你想要的是“第一个流为空的就换第二个”,应该考虑switchIfEmpty,后面实操部分会提到。
1.2 mergeWith:同时订阅与乱序合并
mergeWith和concatWith最大的区别在于:它会同时订阅当前 Flux 和传入的 Publisher,两个流谁先发出数据,就先发射谁的数据。
Flux<Integer> fast = Flux.just(1, 2, 3).delayElements(Duration.ofMillis(100)); Flux<Integer> slow = Flux.just(10, 20, 30).delayElements(Duration.ofMillis(200)); fast.mergeWith(slow) .subscribe(System.out::println);这个例子里,fast的1会先出现,然后是slow的10,接着fast的2……最终的输出顺序是不确定的。所以mergeWith适合那些大体上独立、不要求顺序、只需要把多个信号汇总的场景,比如同时收集多个异步任务的成功/失败结果。
mergeWith内部实际上走的是 flatMap 那一套,对上游所有发布者以并发方式订阅,并用队列缓冲还没被下游请求的数据。如果某个源数据产生太快、另外一个源处理得太慢,合并后的序列就会在内存里积压。这个积压问题我会在第五节专门讲。
1.3 zip:一对一配对的“弱耦合”
zip的语义和前两者完全不同:它不追求“顺序”或者“乱序”,而是做一对一配对。zipWith会将当前 Flux 的元素和另一个 Flux 的元素按位置组合起来,组合的结果由你提供的 BiFunction 决定。
Flux<String> names = Flux.just("Lilei", "HanMeimei"); Flux<Integer> scores = Flux.just(90, 88); names.zipWith(scores, (name, score) -> name + ": " + score) .subscribe(System.out::println); // 输出: // Lilei: 90 // HanMeimei: 88如果两个流的元素数量不一致,多出来的部分会被直接丢弃。zip更高级的用法是Mono.zip,它会把多个Mono同时订阅,等所有Mono都发出一个值之后,再统一回调一次。Mono.zip是聚合接口并行化的核心操作符,到第三节我会实际演示。
三者的对比可以浓缩成一张表,选型的时候对照着看:
| 操作符 | 订阅时间 | 发射语义 | 适合场景 |
|---|---|---|---|
| concatWith | 第一个流结束后再订阅第二个 | 保持先后顺序 | 必须按顺序处理的两个阶段 |
| mergeWith | 同时订阅 | 按到达时间乱序发射 | 独立信号汇总、事件流合并 |
| zip / zipWith | 同时订阅 | 一对一配对,凑不齐就等 | 多路并行结果聚合、特征对齐 |
2. 延迟操作符:delayElements、delaySubscription、delaySequence 与 delayUntil
延迟操作符近几年被问得特别多,因为响应式流里的“延迟”不像Thread.sleep那样直白,它涉及订阅、发射、调度器等多个环节。很多新人以为delayElements就是“让流晚一点输出”,结果一用发现整个流连订阅都比预期晚了,或者每个元素之间的间隔跟预期完全不一样。
2.1 delayElements:每个发射间隔的“呼吸节律”
delayElements(Duration)的作用是:让 Flux 的每个元素在发射前都等待一个固定时长。这个操作符默认把元素切换到Schedulers.parallel()调度器上执行延迟逻辑。
Flux.just("a", "b", "c") .delayElements(Duration.ofSeconds(2)) .subscribe(System.out::println); // 订阅后大约 2 秒输出 a // 再过 2 秒输出 b // 再过 2 秒输出 c这里的重点在于:delayElements改变的是元素与元素之间的时间间隔,它不会把“订阅”本身延后。如果你期望的是“整个流晚 5 秒再开始”,那应该用下一个操作符。
2.2 delaySubscription:延迟的其实是订阅动作
delaySubscription(Duration)从名字上看也带 delay,但它的语义是完全不同的:它延迟的是“订阅上游”这个动作,而不是每个元素的发射。
Flux.just("a", "b", "c") .delaySubscription(Duration.ofSeconds(5)) .subscribe(System.out::println); // 订阅动作会延迟 5 秒才发生 // 一旦订阅完成,a b c 会以极快的速度连续输出这个操作符非常适合用来模拟上游慢的情况,或者在测试里制造“晚到的请求”。我在实际项目中还用它做过客户端冷启动时的局部流量预热:让某些非关键请求延迟几秒开始,等外围服务缓存热起来再真正发出去。
delaySubscription和delayElements几乎是最容易被混淆的一对。可以这样记:delaySubscription是“晚开工”,delayElements是“干活的时候慢一点”。
2.3 delaySequence 与 delayUntil:面向序列与面向信号的延迟
delaySequence(Duration)会让整个序列的所有信号——包括onNext、onComplete、onError——都延后一个固定时长。它相当于把整条时间线平移,但元素之间的相对间隔保持不变。例如:
Flux.just("a", "b", "c") .delaySequence(Duration.ofSeconds(3)) .subscribe(System.out::println); // 订阅后大约 3 秒才开始输出 // a b c 之间的间隔保持原来的样子(这里几乎连续输出)delayUntil则更加灵活:它可以针对每个元素,等待一个由该元素决定的 Publisher 发出信号之后,才发射这个元素。比如一个订单元素进来,先等缓存里的某个配置加载完成,再继续往下游走:
Flux<Order> orders = orderService.loadOrders(); orders.delayUntil(order -> cacheService.warmUp(order.getId())) .subscribe(order -> log.info("emit: {}", order.getId()));这里每个订单都会被delayUntil拦住,直到warmUp返回的 Mono 完成。和delayElements不同,delayUntil的等待时长可以按元素动态变化,灵活性高很多。
3. 聚合接口实操:串行调用改并行的完整过程
介绍完操作符,我拿一个真实的聚合接口案例走一遍。假设现在要实现一个移动端首页接口,需要返回三个部分:用户基本信息、最近订单列表、库存状态。如果不做任何响应式设计,传统写法是一个 Controller 方法里依次调用三个 Feign 接口,耗时是三次调用之和。改造成 WebFlux 之后,很多人也只会用flatMap一层层套,性能几乎没提升,因为本质上还是串行。
3.1 先写一个串行版本,看清性能瓶颈
先把三个下游调用包装成Mono,这是 WebFlux 代码的基础形态:
@GetMapping("/home") public Mono<HomeResponse> home(@RequestParam Long userId) { Mono<UserInfo> userMono = userClient.getUser(userId); Mono<List<OrderInfo>> orderMono = orderClient.getOrders(userId); Mono<StockStatus> stockMono = stockClient.getStock(userId); return userMono.flatMap(user -> orderMono.flatMap(orders -> stockMono.map(stock -> HomeResponse.of(user, orders, stock)))); }这个版本能跑,但三个远程调用是顺序执行的:先查用户,拿到用户之后再查订单,之后才查库存。如果每个下游耗时 200ms,整体就是 600ms 起步。问题在于:三个调用之间根本没有数据依赖,没必要一个等一个。
3.2 用 zip 让三个下游并行执行
这时候Mono.zip的价值就体现出来了。Mono.zip会同时订阅传入的所有 Mono,全部就绪后再组合回调:
@GetMapping("/home") public Mono<HomeResponse> home(@RequestParam Long userId) { Mono<UserInfo> userMono = userClient.getUser(userId); Mono<List<OrderInfo>> orderMono = orderClient.getOrders(userId); Mono<StockStatus> stockMono = stockClient.getStock(userId); return Mono.zip(userMono, orderMono, stockMono) .map(tuple -> HomeResponse.of( tuple.getT1(), tuple.getT2(), tuple.getT3())); }改完之后三个调用同时发起,整体耗时约等于最慢的那个下游,而不是三者之和。假设三个接口都耗时 200ms,串行是 600ms,并行改造后就是 200ms 左右。这是 WebFlux 聚合接口收益最直观的一步。
注意一点:Mono.zip能“并行”是因为它同时订阅了三个上游。如果某个上游内部其实是阻塞调用(比如 JDBC、同步 HTTP 封装),那它依然会占住当前线程。这种阻塞型调用应该包一层subscribeOn(Schedulers.boundedElastic()),把阻塞动作丢到弹性线程池里,否则用zip并行也白搭。
3.3 给聚合加上总超时、降级与空值兜底
聚合接口一旦并行,就会引入新的问题:如果最慢的那一路永远是 10 秒才返回,整个首页接口也跟着 10 秒才出数据。所以必须给整个聚合链路设置超时,超时之后走降级。
@GetMapping("/home") public Mono<HomeResponse> home(@RequestParam Long userId) { Mono<UserInfo> userMono = userClient.getUser(userId) .defaultIfEmpty(UserInfo.EMPTY) .onErrorResume(e -> Mono.just(UserInfo.EMPTY)); Mono<List<OrderInfo>> orderMono = orderClient.getOrders(userId) .onErrorResume(e -> Mono.just(List.of())); Mono<StockStatus> stockMono = stockClient.getStock(userId); return Mono.zip(userMono, orderMono, stockMono) .timeout(Duration.ofSeconds(3)) .onErrorResume(e -> HomeResponse.fallback(ServiceUnavailableError.class)) .map(tuple -> HomeResponse.of( tuple.getT1(), tuple.getT2(), tuple.getT3())); }这里有几个细节值得展开。第一个是defaultIfEmpty,它解决的是上游返回空Mono的情况。Mono如果empty(),zip会一直等不到这个元素,导致整个zip挂住,所以一定要用defaultIfEmpty给它一个兜底值。第二个是timeout,它负责把“最慢的一路”控制在一个可接受范围内。第三个是onErrorResume,它在超时或者其他异常发生的时候返回降级数据。
还有一个常用的组合是switchIfEmpty,它和defaultIfEmpty的区别在于:switchIfEmpty会切换到一个全新的 Publisher,适合“缓存没命中就去查数据库”的场景。简单说,defaultIfEmpty给一个固定默认值,switchIfEmpty给一条新的处理链路。
4. 我踩过的坑:线程切换、上下文丢失与“一慢全慢”
这一节我专门把实际排查过的几个问题完整写出来。如果只给结论不给过程,下次遇到还是会绕远路。
4.1 delayElements 引发的线程切换和 MDC 丢失
现象:网关里用delayElements给某些慢接口做平滑限速,日志突然丢失了全程的 traceId,只有请求入口那一条日志有链路号,后面全变成空白。
排查链路是这样的:先检查了 WebFilter 里打印日志的位置,正常;接着检查 Controller,日志里也还有 traceId;直到往链路深处加了更多日志才发现,从delayElements之后所有日志的 MDC 都空了。最后定位到根因:delayElements内部会把元素调度到Schedulers.parallel()的线程上执行,线程一切换,存在 ThreadLocal 里的 MDC 信息自然就丢了。
这个问题的本质是:Reactor 推崇用Context而不是 ThreadLocal 传递链路数据。但很多日志框架的 MDC 就是 ThreadLocal,Spring Cloud Sleuth/Micrometer Tracing 虽然做了适配,可一旦操作符内部切换了调度器,跨线程传播仍然要靠框架层去做,不一定覆盖所有自定义场景。
如果你在 WebFlux 里想通过 MDC 传 traceId,又不得不用delayElements,一个比较稳的做法是在延迟之前把 traceId 捕获到局部变量,到了延迟之后的阶段再手动写入 MDC:
String traceId = MDC.get("traceId"); Flux.just(...) .delayElements(Duration.ofMillis(500)) .doOnNext(x -> MDC.put("traceId", traceId)) ...现在新项目我更推荐直接用contextWrite把 traceId 写进 Reactor Context,也便于下游从reactor.util.context里取。
4.2 zip 等待最慢源的“木桶效应”
现象:一个聚合接口并行改造后,平均 RT 反而下降了,但 p99 经常飙到 10 秒以上。查监控发现,三个下游里有一个偶尔会非常慢,zip又必须等”全部元素”到达才整体往下走,所以一个慢源拖着整个链路。
这个过程很典型:zip不是“谁先到谁先走”,它是“都到了才一起走”。如果其中一路偶发慢,整个聚合就被这唯一的慢源卡住。我当时的解决思路是给每个子调用单独加timeout,超时后返回默认值,而不是只给zip一个整体超时:
Mono<UserInfo> userMono = userClient.getUser(userId) .timeout(Duration.ofMillis(800)) .onErrorResume(e -> Mono.just(UserInfo.EMPTY));这样受保护的是“每一路”的最长等待时间,即便某一路上游挂掉,zip依然能凑齐数据继续往下走。整体timeout和局部timeout不是替代关系,而是配合关系:局部超时做兜底,整体超时做保险。
4.3 组合操作符的错误传播差异
这也是一个典型的混淆点。concatWith和mergeWith在处理错误时的行为是一致的:任何一个源抛错,错误会直接传播给下游,另一个源的数据即使已经准备好了也不会继续发射。这导致了一个常见的误用场景:有人想用concatWith实现“主数据失败就用备用数据”,但这不管用。
// 错误示例:second 不会因为 first 的 error 而启动 Mono.error(new RuntimeException("boom")) .concatWith(Mono.just("fallback")) .subscribe(...); // 只收到 error正确做法是在第一个流上做错误恢复,或者用onErrorResume切换新的流:
Mono.error(new RuntimeException("boom")) .onErrorResume(e -> Mono.just("fallback")) .subscribe(...); // 正常输出 fallback如果你希望的是“两个源都订阅,谁先出成功结果用谁”,那是firstWithValue或firstWithSignal这类操作符的职责,不要和concatWith混在一起。
5. 从背压和调度器角度看组合与延迟的代价
最后一个部分,聊点容易在压测阶段才暴露的问题。组合操作符和延迟操作符看起来只是改变了数据流的行为,但它们对内存、线程模型和背压的影响,往往比业务逻辑本身更大。
5.1 delayElements 的背压缓冲问题
delayElements的实现逻辑是把元素按延迟节奏“重新产生”一次。上游源通常不会因为你要延迟而停下生产,如果上游生产速度快,下游消费速度被延迟操作拖慢,中间就必然有一个队列在缓冲数据。
我在一个事件推送服务里用delayElements控制推送频率,压测时发现内存占用以肉眼可见的速度上涨。原因就是上游每秒能生产几千个事件,而delayElements设置成了每事件间隔 1 秒,所有来不及发射的事件全部堆积在缓冲队列里。类似场景下,更稳妥的做法是在上游就做限流,或者用limitRate主动控制请求量,而不是靠延迟操作符来“削峰”。
5.2 parallel 与 boundedElastic 调度器如何选
前文提过,delayElements、delaySequence默认都使用Schedulers.parallel(),而delaySubscription默认也是 parallel。parallel线程池的线程数由 CPU 核数决定,适合 CPU 密集型延迟任务;但阻塞型任务(比如远程调用)不应该占用 parallel 线程,应该用Schedulers.boundedElastic()。
如果你写的是网关或聚合层,判断依据很简单:这条链路上有没有真正的阻塞调用。有阻塞,就用subscribeOn(boundedElastic());没有阻塞,保持默认即可。不要为了省事把所有操作都丢给 boundedElastic,因为它的线程数上限有约束,大量堆积的阻塞任务一样能拖垮应用。
5.3 高并发场景下更稳的替代写法
现在回头看,很多“为了演示延迟操作符而用它”的代码,其实可以用更稳的方式替代。
比如想限制某个接口的访问速率,delayElements并不是一个精确的限流器,它只负责把元素铺开,不负责统计单位时间内的请求量。更可靠的限流还是建议用RateLimiter或者limitRate。再比如模拟慢调用,delaySubscription可以做一次性延迟,但如果下游要模拟的是“持续慢响应”,用interval生成时间信号配合take更直观:
Flux.interval(Duration.ofSeconds(1)) .take(3) .flatMap(tick -> someFastMono())这类写法的优势是不会在背压路径上积压太多元素,因为时间信号本身就控制了发射节奏。
组合和延迟操作符本身不复杂,复杂的是它们和你系统的线程模型、背压策略纠缠在一起后产生的各种隐性代价。我在项目里总结了一条经验:凡是加延迟操作符的地方,都要想一想“如果上游产量超过下游处理能力,数据会堆在哪”。想清楚这个问题,很多线上事故在上线前就能避免。