请解释 Flux.merge() 的概念,并说明在 AI 大模型评测平台中,如何利用该方法对多个模型进行并行流式调用?
考察说明
考查对响应式编程中合并操作符的理解及其在真实业务场景(多模型并行调用)中的应用能力。
回答思路
- 【回答框架 1】Flux.merge() 是 Project Reactor 中的静态方法,用于将多个 Flux 或 Mono 流合并为一个 Flux,所有源流同时被订阅,发出的元素按到达顺序交织,不保证顺序。
- 【回答框架 2】与 concat 不同,merge 是并发订阅,适合并行执行耗时操作;若需要按顺序处理,应使用 concat。
- 【回答框架 3】在 AI 大模型评测平台中,可用 Flux.merge(List<Flux<ModelResult>>) 同时向多个模型发起请求,每个请求返回 Flux,合并后统一订阅,实现流式输出。
- 【回答框架 4】需注意背压和错误处理:merge 默认错误传播可能导致整体中断,可通过 onErrorContinue 或 onErrorResume 单独处理每个流。
- 【回答框架 5】可使用 Scheduler 控制并发度,如 Schedulers.parallel(),并根据平台资源设置合适的并发数,避免过载。
- 【关键点 1】Flux.merge() 并发订阅多个流,按到达顺序合并。
- 【关键点 2】与 concat 的区别:merge 无序并发,concat 有序串行。
- 【关键点 3】使用 onErrorContinue 或 onErrorResume 实现单个流错误隔离。
- 【关键点 4】通过 Scheduler 限制并发度,保护下游系统。
- 【易错点 1】误以为 merge 保证顺序,实际不保证。
- 【易错点 2】背压处理不当可能导致内存压力。
- 【易错点 3】全局错误处理可能导致一个模型失败影响整体。