为何基于IntStream实现的自定义limit收集器无法让流短路?
为何基于IntStream实现的自定义limit收集器无法让流短路?
嘿,我前阵子刚在Java 24里捣鼓刚转正的Stream Gatherers,想自己复刻个limit操作练练手,结果踩了个大坑——用在IntStream上的时候,明明已经取够指定数量的元素了,后面的元素居然还在被处理,完全不短路!先给你看看我当时写的半截代码:
public class LimitDemo { private static class State { int count; } static <T> Gatherer<T, ?, T> limit(int size) { Supplier<State> initializer = State::new; Gatherer.Integrator<State, T, T> integrator = (state, element, downstream) -> { if (state.count < size) { downstream.push(element); state.count++; return true; } return false; }; // 当时我没加关键的特征声明! return Gatherer.of(initializer, integrator); } }
后来查了半天终于搞懂问题出在哪了,给你唠唠核心原因:
没给Gatherer声明短路特征:流框架可不是光看你的Integrator返回false就会终止的,它得先知道这个Gatherer支持短路啊!如果自定义Gatherer没有显式声明
SHORT_CIRCUIT特征,流框架会默认认为这个操作必须处理完所有元素,哪怕你的逻辑里已经不想继续了。原始流适配的隐形坑:我一开始写的是通用类型的Gatherer,套在IntStream上用的时候,底层会自动把int装箱成Integer,这个适配过程会多一层包装,有时候会把短路的信号给“吞”了。最好是针对原始流写专门的
Gatherer.OfInt实现,直接处理原始类型,信号传递更直接。
那怎么改呢?给你看修正后的代码:
首先是通用版加短路特征声明:
public class LimitDemo { private static class State { int count; } static <T> Gatherer<T, ?, T> limit(int size) { Supplier<State> initializer = State::new; Gatherer.Integrator<State, T, T> integrator = (state, element, downstream) -> { if (state.count < size) { downstream.push(element); state.count++; return true; } return false; }; // 关键!加上短路特征声明,告诉流框架这个操作可以提前终止 return Gatherer.of(initializer, integrator) .withCharacteristics(Gatherer.Characteristics.SHORT_CIRCUIT); } }
然后是专门给IntStream用的原始类型版本,效率更高还不容易出问题:
public class LimitDemo { private static class IntState { int count; } static Gatherer.OfInt limitInt(int size) { Supplier<IntState> initializer = IntState::new; Gatherer.OfInt.Integrator<IntState> integrator = (state, element, downstream) -> { if (state.count < size) { downstream.push(element); state.count++; return true; } return false; }; return Gatherer.ofInt(initializer, integrator) .withCharacteristics(Gatherer.Characteristics.SHORT_CIRCUIT); } }
这么一改之后,IntStream用这个自定义limit就会正常短路了——取够指定数量的元素后,后面的元素直接跳过,不会再被处理。
内容来源于stack exchange
相关产品推荐
相关产品推荐

