Java Gatherer链式调用时Downstream.isRejecting状态不一致问题
链式Gatherer中isRejecting状态未正确更新的问题分析
问题现象
当链式调用两个Gatherer时:
gather2的Integrator返回false(表示拒绝接收更多元素)gather1的finisher中调用downstream.push()返回false,这符合预期- 但
downstream.isRejecting()始终返回false,与下游实际的拒绝状态不符
复现代码
Stream.of(0) .gather(gather1()) .gather(gather2()) .forEach(System.out::println); static Gatherer<Integer, ?, Integer> gather1() { return Gatherer.ofSequential( Gatherer.Integrator.ofGreedy((state, element, downstream) -> downstream.push(element)), (unused, downstream) -> { for (int i = 1; i <= 10 && !downstream.isRejecting(); i++) { System.out.println("finisher pushing " + i + ", push result " + downstream.push(i) + ", isRejecting " + downstream.isRejecting()); } } ); } static Gatherer<Integer, ?, Integer> gather2() { return Gatherer.ofSequential( Gatherer.Integrator.of( (state, element, downstream) -> false ) ); }
原因解析
链式调用的Stream中,gather1的下游并非直接指向gather2,而是Stream框架内部的中间包装节点。当下游gather2返回拒绝状态后:
push()方法会向下传递调用并返回下游的拒绝结果,所以能得到false- 但中间包装节点的
isRejecting()方法没有同步更新上游可见的状态,导致上游读取到的始终是初始的false - 这是Gatherer链式实现中的状态同步设计问题,仅
push()的返回值能实时反映当下的下游状态
解决方案
放弃依赖isRejecting()判断下游状态,直接使用push()的返回值作为终止条件:
static Gatherer<Integer, ?, Integer> gather1() { return Gatherer.ofSequential( Gatherer.Integrator.ofGreedy((state, element, downstream) -> downstream.push(element)), (unused, downstream) -> { for (int i = 1; i <= 10; i++) { boolean pushSuccess = downstream.push(i); System.out.println("finisher pushing " + i + ", push result " + pushSuccess + ", isRejecting " + downstream.isRejecting()); if (!pushSuccess) { break; } } } ); }
内容的提问来源于stack exchange,提问作者RamPrakash
相关产品推荐
相关产品推荐

