Spring Reactor:如何按Key等待多个Flux处理完成?
嘿,我明白你的问题了——你有个无限发射IP的数据源,想用两个处理器分别处理每个IP,然后把同一个IP的两个结果合并传给接收器,但原代码因为流永远不会结束,导致collectMap一直没法输出结果,对吧?
问题根源
你原代码里的collectMap是一个终止操作,它需要等待整个流结束才能收集完所有元素并返回结果。但你的数据源用了repeat()变成了无限流,每个IP分组对应的Flux会不断产生新元素,永远不会终止,所以collectMap永远不会触发后续的打印逻辑。
解决方案
根据你的场景,我给你两种实用的实现方式:
方案1:用zip配对(适用于结果顺序严格对应的场景)
因为你用了publish()让两个处理器共享同一个数据源,每个IP会同时被两个map处理,生成的Foo在两个Flux里是严格一一对应的。这种情况下直接用zip配对最简洁高效:
public class Demo { public static void main(String[] args) throws Exception { Flux<String> source = Flux.fromIterable(Lists.newArrayList("1.1.1.1", "2.2.2.2", "3.3.3.3")) .delayElements(Duration.ofMillis(500)) .repeat(); ConnectableFlux<String> ipsFlux = source.publish(); Flux<Foo> fooFlux1 = ipsFlux.map(ip -> new Foo(ip, "1")); Flux<Foo> fooFlux2 = ipsFlux.map(ip -> new Foo(ip, "2")); // 直接配对两个Flux中位置对应的元素(同一个IP的两个处理结果) Flux.zip(fooFlux1, fooFlux2) .map(pair -> { Foo foo1 = pair.getT1(); Foo foo2 = pair.getT2(); // 合并成你需要的Map结构 Map<String, String> result = new HashMap<>(); result.put(foo1.type, foo1.id); result.put(foo2.type, foo2.id); return result; }) .subscribe(System.out::println); ipsFlux.connect(); Thread.currentThread().join(); } static class Foo { String id; String type; public Foo(String id, String type) { this.id = id; this.type = type; } public String getId() { return id; } @Override public String toString() { return "Foo{" + "id='" + id + '\'' + ", value='" + type + '\'' + '}'; } } }
方案2:用groupBy+take(2)(适用于结果可能乱序的场景)
如果两个处理器的处理速度不一致,导致同一个IP的两个Foo到达顺序混乱,你可以用分组+取固定数量元素的方式,确保收集到同一个IP的两个结果就立即处理:
public class Demo { public static void main(String[] args) throws Exception { Flux<String> source = Flux.fromIterable(Lists.newArrayList("1.1.1.1", "2.2.2.2", "3.3.3.3")) .delayElements(Duration.ofMillis(500)) .repeat(); ConnectableFlux<String> ipsFlux = source.publish(); Flux<Foo> fooFlux1 = ipsFlux.map(ip -> new Foo(ip, "1")); Flux<Foo> fooFlux2 = ipsFlux.map(ip -> new Foo(ip, "2")); Flux.merge(fooFlux1, fooFlux2) .groupBy(Foo::getId) .flatMap(group -> // 每个IP分组取2个元素(对应两个处理器的结果),收集完成后立即返回 group.take(2) .collectMap(Foo::getType, Foo::getId) ) .subscribe(System.out::println); ipsFlux.connect(); Thread.currentThread().join(); } // Foo类同上 }
核心思路
两种方案都是避免等待整个无限流终止,而是每凑齐一组对应结果就立即处理:
- 方案1利用共享数据源的顺序一致性,直接配对元素;
- 方案2通过分组+限定元素数量,确保每个IP的两个结果收集完成就触发处理。
内容的提问来源于stack exchange,提问作者IsaacLevon
相关产品推荐
相关产品推荐

