Reactor groupJoin方法TrightEnd为UnicastProcessor且size=0问题咨询
Reactor groupJoin 无输出原因及工作机制解析
问题根源:参数理解错误
你写的groupJoin代码里,第四个参数的y是Flux类型(不是单个元素),但你直接把它当成单个元素拼接,而且没有订阅这个Flux,自然不会有任何元素被触发发射,所以看不到输出。
而join的第四个参数中,y是右侧流的单个元素,所以直接拼接就能正常输出所有配对组合。
groupJoin 工作机制
groupJoin的核心是为左侧流的每个元素,配对右侧流中所有在该左元素生命周期内发射的元素组成的Flux,具体规则:
- 当左侧流发射一个元素
x时,启动leftEnd(第二个参数)定义的生命周期信号流;只要这个信号流没完成,x就处于"活跃"状态 - 右侧流发射元素
y时,启动rightEnd(第三个参数)定义的生命周期信号流;只要这个信号流没完成,y就处于"活跃"状态 - 所有处于活跃状态的
x,会和所有处于活跃状态的y所属的Flux进行配对 - 最终的结果由第四个参数
(x, yFlux) -> ...生成,这里的yFlux是所有符合条件的右侧元素的集合流,必须订阅它才能获取元素
修正后的代码(实现你期望的分组效果)
把第四个参数改成收集右侧Flux的元素,比如用collectList()来获取所有配对的右侧元素:
public static void testGroupJoin(){ Flux<Integer> f1 = Flux.just(1,2,3,10,11,12,13,14); Flux<Integer> f2 = Flux.just(10,12,13,14,15,16); f1.groupJoin(f2, x -> Flux.never(), // 左元素永远活跃 y -> Flux.never(), // 右元素永远活跃 (x, yFlux) -> yFlux.collectList().map(list -> x + ":" + list) ).subscribe(System.out::println); }
运行后会输出:
1:[10,12,13,14,15,16] 2:[10,12,13,14,15,16] 3:[10,12,13,14,15,16] 10:[10,12,13,14,15,16] 11:[10,12,13,14,15,16] 12:[10,12,13,14,15,16] 13:[10,12,13,14,15,16] 14:[10,12,13,14,15,16]
groupJoin vs join 的核心差异
join:每一对符合活跃条件的左元素+右元素,单独生成一个结果(比如你的join代码会输出1,10、1,12、2,10等所有组合)groupJoin:每个左元素,和所有符合条件的右元素组成的Flux生成一个结果(即按左元素分组,把所有匹配的右元素归为一组)
内容的提问来源于stack exchange,提问作者Mr.Wang
相关产品推荐
相关产品推荐

