如何不使用subscribe()基于Mono<List<String>>迭代值条件调用函数
问题解答
问题2 同一个Mono是否支持多次订阅
Mono默认是冷流,允许多次调用subscribe(),但每次订阅都会触发上游整条链路的重新执行:你在checkCondition()里订阅一次、buildResponse()的zip操作隐式订阅一次,等于两次独立运行了merge()里的全量逻辑,所以doOnSuccess()会被触发两次,属于正常特性,并非Mono只能订阅一次。
问题1 避免doOnSuccess重复触发的解决方案
有三种常见实现方案,可根据业务场景选择:
方案1:对公共Mono增加缓存(最推荐,改动最小)
在merge()返回的Mono后添加.cache()算子,上游执行一次后的结果会被缓存,后续所有订阅都会复用同一份结果,不会重新触发上游逻辑,doOnSuccess()自然只会执行一次。
修改后的merge()代码如下:
private fun merge(list1: Mono<List<String>>, list2: Mono<List<String>>) = Flux.merge( list1.flatMapMany { Flux.fromIterable(it) }, list2.flatMapMany { Flux.fromIterable(it) } ) .collectList() .cache() // 新增缓存逻辑,多订阅复用结果 .doOnSuccess { LOG.debug("List of words: $it") }
方案2:调整doOnSuccess的挂载位置
如果不需要每次获取listOfWords都打日志,只需要在构造响应时打印,就把doOnSuccess()从公共的merge()方法里移除,改到buildResponse()的zip操作之后挂载,这样checkCondition()里的订阅不会触发日志打印。
方案3:把校验逻辑合并到响应流,避免提前订阅
你当前在checkCondition()里主动调用subscribe()属于脱离响应式流的独立订阅,不符合声明式编程规范,可以把校验逻辑改为非订阅的流式写法,组合到最终的响应流中,全程只有一次订阅,从根源上避免重复执行的问题。
修改后的checkCondition()代码如下:
private fun checkCondition( listOfWords: Mono<List<String>>, ): Mono<Void> { return listOfWords.doOnNext { it.forEach { word -> if (someCondition(word)) { alarmSystem.notify("Something is missing for word {0}") } } }.then() }
之后在buildResponse()中把校验逻辑和现有zip逻辑组合即可。
内容的提问来源于stack exchange,提问作者marcin.winny
相关产品推荐
相关产品推荐

