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

如何合并两个含Pair的Flux流,按共享元素匹配并补全缺失值

按键合并Flux流的简洁高效实现方案

这是个很常见的按键合并反应式流的需求,刚好Reactor提供了简洁高效的实现方式,分两种匹配场景给你说明:

场景1:保留s1的所有键,匹配s2中存在的键(完全符合你的示例需求)

这种情况最直接高效的方式是先把s2收集为一个内存级的键值对Map,再遍历s1的每个元素去匹配对应的值:

import reactor.core.publisher.Flux;
import reactor.core.publisher.Mono;
import reactor.util.function.Tuples;
import java.util.Map;

// 假设你的流元素是Tuple2<Integer, String>,如果是自定义键值类,替换成对应的getKey/getValue方法即可
Flux<Tuple2<Integer, String>> s1 = Flux.just(Tuples.of(1, "A"), Tuples.of(2, "B"), Tuples.of(3, "C"));
Flux<Tuple2<Integer, String>> s2 = Flux.just(Tuples.of(1, "X"), Tuples.of(3, "XX"));

// 异步收集s2的键值对为Map,仅执行一次
Mono<Map<Integer, String>> s2KeyMap = s2.collectMap(Tuple2::getT1, Tuple2::getT2);

// 遍历s1,结合s2的Map生成配对结果
Flux<Tuple2<Integer, Tuple2<String, String>>> result = s1.flatMap(s1Element ->
    s2KeyMap.map(map ->
        // 生成Pair,s2中无对应键则用null填充
        Tuples.of(s1Element.getT1(),
            Tuples.of(s1Element.getT2(), map.getOrDefault(s1Element.getT1(), null))
        )
    )
);

// 验证输出结果
result.subscribe(System.out::println);

执行后输出结果完全符合你的期望:

[1, [A,X]]
[2, [B,null]]
[3, [C,XX]]

这种实现的优势是高效且简洁:仅需一次收集s2的操作,遍历s1时直接在内存中查询匹配,时间复杂度为O(n+m)(n为s1元素数,m为s2元素数)。

场景2:保留s1和s2的所有键(键的并集)

如果你的需求扩展为需要包含两个流中所有的键(比如s2存在s1没有的键也要保留),可以先收集所有唯一键,再分别匹配两个流的值:

import reactor.core.publisher.Flux;
import reactor.core.publisher.Mono;
import reactor.util.function.Tuples;
import java.util.Map;
import java.util.Set;
import java.util.stream.Collectors;

Flux<Tuple2<Integer, String>> s1 = Flux.just(Tuples.of(1, "A"), Tuples.of(2, "B"), Tuples.of(3, "C"));
Flux<Tuple2<Integer, String>> s2 = Flux.just(Tuples.of(1, "X"), Tuples.of(3, "XX"), Tuples.of(4, "Y"));

// 收集两个流的所有唯一键
Mono<Set<Integer>> allKeys = Flux.merge(s1.map(Tuple2::getT1), s2.map(Tuple2::getT1))
    .collect(Collectors.toSet());

// 分别收集两个流的键值Map
Mono<Map<Integer, String>> s1Map = s1.collectMap(Tuple2::getT1, Tuple2::getT2);
Mono<Map<Integer, String>> s2Map = s2.collectMap(Tuple2::getT1, Tuple2::getT2);

// 遍历所有键,匹配两个流对应的值
Flux<Tuple2<Integer, Tuple2<String, String>>> result = allKeys.flatMapMany(keys ->
    Flux.fromIterable(keys)
        .flatMap(key ->
            Mono.zip(
                s1Map.map(map -> map.get(key)), // s1中无对应键则返回null
                s2Map.map(map -> map.get(key))  // s2中无对应键则返回null
            ).map(pair -> Tuples.of(key, pair))
        )
);

result.subscribe(System.out::println);

执行后会包含s2中独有的键4:

[1, [A,X]]
[2, [B,null]]
[3, [C,XX]]
[4, [null,Y]]

注意事项

  • 以上实现基于有限流:如果你的流是无限流,collectMap会一直等待流结束,无法正常工作。这种情况下需要用基于窗口或分组的反应式处理(比如groupBy结合bufferTimeout),但复杂度会有所提升。
  • 如果你的元素不是Tuple2,可以替换成自定义的键值对类,只要提供获取键和值的方法即可。

内容的提问来源于stack exchange,提问作者Pedro Alipio

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.28 07:28:30