ARTICLE DETAIL

资讯详情

深耕网站视觉设计与运营推广的一线实战洞察。

响应式编程必读:Flux与Mono组合操作符及延迟操作符实战解析

响应式编程必读:Flux与Mono组合操作符及延迟操作符实战解析 写响应式代码的人大概都有过这种体验Flux 和 Mono 表面上就是把数据包成一条流可真到了要把两条流合在一起、或者想让流稍微“慢半拍”的时候组合操作符和延迟操作符的选择立刻变得不直观了。我在一个 Spring WebFlux 网关项目里就因为在“等前一批跑完”和“谁先到谁先走”之间选错了操作符上线第一天订单推送顺序乱掉被迫紧急热修复。后来我把这些操作符的底层行为、调度器选择和适用场景彻底过了一遍这篇就当作一次完整的复盘。这篇内容主要讲两块一是组合操作符concatWith、mergeWith、zip 以及它们的变体二是延迟操作符delayElements、delaySubscription、delaySequence、delayUntil最后落到实际聚合接口的并行化改造上。无论你是在写聚合 API、做网关转发还是刚接触响应式想少踩几个坑这篇都值得看完。1. 组合操作符concatWith、mergeWith、zip 的底层行为差异很多人第一眼看到这三个操作符会觉得它们做的事情都一样把两个流合并成一个。实际上它们的差异非常大用错之后的表现也不一样——有的是顺序错乱有的是数据丢失有的是整个流被最慢的一路“锁死”。我自己的经验是先分清三者的语义模型再记用法否则代码写对了也是一头雾水。1.1 concatWith接力棒式的顺序连接concatWith的语义是先订阅当前 Flux等它完整结束后再订阅传入的另一个 Publisher然后把第二个流的数据接在第一个流后面。FluxInteger first Flux.just(1, 2, 3); FluxInteger 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两个流谁先发出数据就先发射谁的数据。FluxInteger fast Flux.just(1, 2, 3).delayElements(Duration.ofMillis(100)); FluxInteger 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 决定。FluxString names Flux.just(Lilei, HanMeimei); FluxInteger 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 发出信号之后才发射这个元素。比如一个订单元素进来先等缓存里的某个配置加载完成再继续往下游走FluxOrder 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 MonoHomeResponse home(RequestParam Long userId) { MonoUserInfo userMono userClient.getUser(userId); MonoListOrderInfo orderMono orderClient.getOrders(userId); MonoStockStatus 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 MonoHomeResponse home(RequestParam Long userId) { MonoUserInfo userMono userClient.getUser(userId); MonoListOrderInfo orderMono orderClient.getOrders(userId); MonoStockStatus 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 MonoHomeResponse home(RequestParam Long userId) { MonoUserInfo userMono userClient.getUser(userId) .defaultIfEmpty(UserInfo.EMPTY) .onErrorResume(e - Mono.just(UserInfo.EMPTY)); MonoListOrderInfo orderMono orderClient.getOrders(userId) .onErrorResume(e - Mono.just(List.of())); MonoStockStatus 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 就是 ThreadLocalSpring Cloud Sleuth/Micrometer Tracing 虽然做了适配可一旦操作符内部切换了调度器跨线程传播仍然要靠框架层去做不一定覆盖所有自定义场景。如果你在 WebFlux 里想通过 MDC 传 traceId又不得不用delayElements一个比较稳的做法是在延迟之前把 traceId 捕获到局部变量到了延迟之后的阶段再手动写入 MDCString 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一个整体超时MonoUserInfo 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())这类写法的优势是不会在背压路径上积压太多元素因为时间信号本身就控制了发射节奏。组合和延迟操作符本身不复杂复杂的是它们和你系统的线程模型、背压策略纠缠在一起后产生的各种隐性代价。我在项目里总结了一条经验凡是加延迟操作符的地方都要想一想“如果上游产量超过下游处理能力数据会堆在哪”。想清楚这个问题很多线上事故在上线前就能避免。
返回列表