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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.13 09:03:46