如何合并两个含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
相关产品推荐
相关产品推荐

