如何从同一输入Flux生成含两个映射对象的目标Flux流?
优化Flux转换方案:避免重复读取源流
你的问题核心是要将源流中每个List<ProgressEvent>转换为两个ServerSentEvent<String>,但当前用mergeWith的方式会导致源流被订阅两次,这不仅低效,还可能引发重复消费或线程安全问题(若源流不是热流)。
最优方案是使用flatMap操作符,它可以将每个输入元素转换为包含多个输出元素的Flux,再将所有小流合并成一个大流,全程仅需订阅源流一次。
示例代码实现
Flux<List<ProgressEvent>> source = ... Flux<ServerSentEvent<String>> target = source .flatMap(events -> { // 生成第一个基于message字段的ServerSentEvent ServerSentEvent<String> messageEvent = ServerSentEvent.<String>builder() .data(events.stream().map(ProgressEvent::getMessage).collect(Collectors.joining("; "))) .event("progress-messages") .build(); // 生成第二个基于value字段的ServerSentEvent double totalProgress = events.stream().mapToDouble(ProgressEvent::getValue).sum(); ServerSentEvent<String> valueEvent = ServerSentEvent.<String>builder() .data(String.format("Total progress: %.2f", totalProgress)) .event("progress-values") .build(); // 将两个事件包装为Flux返回,保证顺序 return Flux.just(messageEvent, valueEvent); });
方案优势
- 仅订阅源流一次:
flatMap逐个处理源流元素,每个List<ProgressEvent>仅被消费一次,避免重复读取带来的性能损耗。 - 逻辑内聚:两个事件的生成逻辑集中在同一个回调中,代码更紧凑易维护。
- 顺序可控:通过
Flux.just()可以严格保证每个输入元素对应的两个事件的输出顺序,而mergeWith的顺序可能受订阅时机影响。
如果需要全局严格的事件顺序,也可以替换为concatMap(串行处理,性能略低于flatMap,按需选择即可)。
内容的提问来源于stack exchange,提问作者Wim Deblauwe
相关产品推荐
相关产品推荐

