Reactor多Flux/Mono编排问题:异步前置执行与动态合并需求
Reactor异步编排:动态Flux管理与提前执行优化
核心业务场景
你说的这个模块要处理实体ID和「解析类型」参数,靠多个Flux异步操作攒数据,核心需求其实就是这几点:
- 部分解析类型得先跑同步前置操作,等同步数据就绪了再启动剩下的Flux
- 想提前启动部分Flux来缩短总耗时,还得保证数据一点都不能丢
- 要能动态管理一堆Flux:支持初始为空、随时加新的、最后合并成一个结果
针对你提的三个问题的具体解法
1. 创建可初始为空的Flux容器
要搞一个初始为空、之后能动态加数据的Flux容器,Reactor 3.2之后推荐用**Sinks.Many**(替代了旧的Processor API)。它就像一个可以随时往里面灌数据的管道,初始的时候啥都没有,完全等价于Flux.empty():
// 搞一个支持多订阅、带背压缓冲的Sink,适合大部分场景 val dataSink = Sinks.many().multicast().onBackpressureBuffer<String>() // 对应的Flux容器,初始是空的,后续所有加进来的数据都会走这个Flux val combinedFlux: Flux<String> = dataSink.asFlux()
如果你的场景是只有一个订阅者,那用Sinks.many().unicast().onBackpressureBuffer()性能会更好。
2. 动态添加Flux并合并为Mono
有了Sinks.Many之后,你随时可以把任意有限异步Flux的数据转发到这个Sink里,不用管数据从哪来。最后只要对combinedFlux调用collectList(),就能得到所有数据合并后的Mono<List<T>>:
// 示例:把一个新的Flux加到容器里 fun addFluxToSink(newFlux: Flux<String>) { newFlux.doOnNext { // 把Flux的每个数据发送到Sink里 dataSink.emitNext(it, Sinks.EmitFailureHandler.FAIL_FAST) } .doOnComplete { // 可选:如果所有要加的Flux都跑完了,可以标记Sink结束 // dataSink.emitComplete() } .subscribe() // 调用subscribe就会立刻启动这个Flux的执行,实现提前启动 } // 最终合并所有数据的Mono val resultMono: Mono<List<String>> = combinedFlux.collectList()
3. 提前启动Flux不丢数据,动态追加同类型Flux
- 提前启动不丢数据:只要在启动Flux之前,已经有下游订阅了
combinedFlux(或者说Sink对应的Flux),Flux产生的数据会被Sink的背压缓冲区存起来,绝对不会丢。就算还没下游订阅,Sink也会默认缓存数据(缓冲区大小默认无界,也可以自己配置)。 - 动态追加同类型Flux:直接给已启动的Flux加数据是不行的,但通过
Sinks.Many当中间层就好办了——你随时可以加新的同类型Flux,把它的数据转发到Sink里,这样combinedFlux就会包含所有后续加进来的Flux的数据,完美实现“维护单个Flux引用还能动态加”的需求。
要是你不想用Sink,也可以用MutableList<Flux<T>>存所有要合并的Flux,最后调用Flux.merge(list)合并,但这种方式没法提前启动部分Flux(必须等所有Flux都加进列表再合并订阅),所以Sink方案更贴合你“提前启动压缩耗时”的需求。
结合你的业务场景的完整Kotlin实现
下面是优化后的实现(保留了你原有的核心逻辑,还加了动态Flux管理的能力,能体现提前启动的优势):
private val log = KotlinLogging.logger {} class ReactiveDataService { // 动态Flux容器:用Sink当中间层,支持随时加新的Flux private val dataSink = Sinks.many().multicast().onBackpressureBuffer<String>() private val combinedFlux = dataSink.asFlux() private val createMono: () -> Mono<List<Int>> = { Flux.just(9, 8, 7) .flatMap { Flux.fromIterable(List(it) { Random.nextInt(0, 100) }) .parallel() .runOn(Schedulers.boundedElastic()) } .collectList() .cache() // 缓存结果,避免重复执行耗时操作 } private val processResults: (List<String>, List<String>) -> String = { d1, d2 -> "\n\tdownstream 1: $d1\n\tdownstream 2: $d2" } private val convert: (List<Int>, Int) -> Flux<String> = { data, multiplier -> Flux.fromIterable(data.map { String.format("%3d", it * multiplier) }) } fun doQuery(): String? { val mono = createMono() // 提前启动第一个下游Flux(不用等同步操作,mono是缓存的,不会重复跑) val downstream1 = mono.flatMapMany { convert(it, 1) } downstream1.doOnNext { dataSink.emitNext(it, Sinks.EmitFailureHandler.FAIL_FAST) } .doOnComplete { log.info("Downstream 1 finished collecting data") } .subscribe() // 模拟同步前置操作(比如解析类型校验、参数准备) Thread.sleep(50) // 模拟同步操作的耗时 // 同步操作完成后,启动第二个下游Flux val downstream2 = mono.flatMapMany { convert(it, 2) } downstream2.doOnNext { dataSink.emitNext(it, Sinks.EmitFailureHandler.FAIL_FAST) } .doOnComplete { log.info("Downstream 2 finished collecting data") dataSink.emitComplete() // 所有Flux都跑完了,标记Sink结束 } .subscribe() // 这里保留你原有的zip逻辑,展示提前启动的效果:两个下游并行执行,总耗时不会是两者之和 val resultMono = Mono.zip( downstream1.collectList(), downstream2.collectList(), processResults ) return resultMono.block() } } fun main() { val service = ReactiveDataService() val start = System.currentTimeMillis() val result = service.doQuery() log.info("{}\n\tTotal time: {}ms", result, System.currentTimeMillis() - start) }
输出结果示例
downstream 1: [ 66, 39, 40, 88, 97, 35, 70, 91, 27, 12, 84, 37, 35, 15, 45, 27, 85, 22, 55, 89, 81, 21, 43, 62] downstream 2: [132, 78, 80, 176, 194, 70, 140, 182, 54, 24, 168, 74, 70, 30, 90, 54, 170, 44, 110, 178, 162, 42, 86, 124] Total time: 209ms
内容的提问来源于stack exchange,提问作者Steve Storck
相关产品推荐
相关产品推荐

