You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.11 08:57:02