Flux.merge 是 Project Reactor 的合并算子,它会同时订阅多个 Flux 源、把发出的元素按时间交错合并到一个流。在 AI 大模型评测平台里,需要同时对多家模型发起流式调用再合并结果,Flux.merge 正好干这个——多路并行、交错输出,比顺序串行快得多。理解它的订阅语义和错误处理,是用好这个算子的关键。
一、Flux.merge 的语义
| 对比项 | merge | concat | flatMap |
|---|---|---|---|
| 订阅时机 | 同时订阅所有源 | 逐个订阅,前一个完成才下一个 | 同时订阅(映射后合并) |
| 输出顺序 | 按到达时间交错 | 严格按源顺序 | 交错 |
| 适用场景 | 多路并行、不关心顺序 | 有序串行 | 一对多映射 |
| 背压 | 合并流各自背压 | 逐个 | 各子流背压 |
merge 的关键特性是「同时订阅」——所有源在订阅阶段就并发启动,谁先出结果谁先进入合并流。评测平台要的就是这种「多模型并行、谁快谁先返回」的行为。concat 会等第一个模型完全结束再开始第二个,完全失去并行性,多模型评测会慢几倍。flatMap 也能并行,但语义是「映射后合并」,需要先映射再订阅,多一层抽象;merge 直接面向「已有多个流、合并它们」的场景,更直白。
二、评测平台里怎么用它做多模型并行流式调用
- 把每个被评测模型包成一个返回
Flux<String>的流式调用; - 用
Flux.merge()合并多个模型的流; - 给合并流加超时和错误兜底,单模型失败不影响其他路;
- 前端 SSE 订阅这个合并流,边出边显示。
@Service
public class ParallelEvalService {
private final Map<String, ChatModel> models; // 多家模型
public Flux<EvalChunk> runParallel(String prompt, Duration timeout) {
List<Flux<EvalChunk>> streams = models.entrySet().stream()
.map(e -> streamFor(e.getKey(), e.getValue(), prompt))
.toList();
return Flux.merge(streams)
.timeout(timeout)
.onErrorResume(e -> Flux.empty()); // 单路失败兜底
}
private Flux<EvalChunk> streamFor(String name, ChatModel model, String prompt) {
return ChatClient.create(model).prompt().user(prompt).stream().content()
.map(chunk -> new EvalChunk(name, chunk));
}
}
这段代码的核心是 Flux.merge(streams)——它同时订阅所有模型的流式输出,谁先吐字谁先进合并流。EvalChunk 带模型名,前端收到 chunk 后按模型名分发到对应的回答区域,实现「多家模型同时打字」的效果。timeout 给整个合并流设上限,防止单个慢模型拖住整体。
三、为什么是 merge 不是 concat 或 flatMap
- concat 会等第一个模型完全结束再开始第二个,如果模型 A 要 30 秒、模型 B 要 5 秒,用 concat 总耗时 35 秒,用 merge 只要 30 秒(B 在 A 跑的同时就跑完了);
- flatMap 也能并行,但它的语义是「把上游每个元素映射成新流再合并」,多了一层映射关系;merge 直接面向「已有 N 个独立流,合并成一个」,场景更贴切;
- 评测平台的多模型调用是「N 个独立流合并」,不需要映射,merge 语义刚好对上。
四、错误处理要注意的点
merge 默认是「一源报错就传播错误」,这在多模型评测里不合适——一个模型挂了不能让整个评测失败。所以必须配 onErrorResume 把单路错误吞掉转成空流,最后用 Flux.collectList 看哪些路成功了。如果想精确知道哪路失败,可以在 onErrorResume 里发一个带模型名的错误 chunk,前端显示「模型 X 推理失败」,其他模型继续展示。
五、背压和流速控制
多家模型并行流式输出时,合并流的数据流速可能远超前端消费速度(比如 5 家模型每秒各吐 50 个 token,合并后每秒 250 个 token)。如果前端 SSE 跟不上,会在服务端积压。解法是给合并流加 onBackpressureBuffer(timeline) 或 limitRate,控制下游消费速率,避免内存爆。评测平台一般用 buffer 模式(容忍短暂积压),因为评测结果不能丢。
常见问题
Q1:merge 会保序吗?
不会。按元素到达时间交错,谁先返回谁在前。
Q2:单路超时怎么处理?
给整路加 timeout,超时后 onErrorResume 吞掉,不影响其他路。
Q3:merge 和 mergeSequential 区别?
merge 是纯交错,mergeSequential 会缓冲后按源顺序输出,延迟更高。
Q4:并行调几家模型合适?
看上游模型 API 并发额度。一般 3~5 家并行,再多容易撞各家限流。