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

如何从同一输入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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.09 10:42:40