Java中使用parallelStream配合reduce出现结果异常且随机的技术问题
问题原因分析
咱们先拆解下这段并行流代码为啥会出现随机重复的异常结果,核心问题出在错误使用reduce方法且违反了并行流reduce的核心契约:
- 共享可变的identity实例导致线程安全与重复合并
你用的reduce(U identity, BiFunction<U,? super T,U> accumulator, BinaryOperator<U> combiner)重载中,第一个参数identity要求是「恒等初始值」——简单说就是,并行流的每个工作线程都应该拿到独立的初始容器,而且这个初始值不能被修改。但你直接传入了new ArrayList<>(),这意味着所有并行线程都会操作同一个ArrayList实例。
看你的执行日志就能发现:不管是accumulator阶段还是combiner阶段,操作的都是同一个列表对象。比如第一次执行时,combiner把[b1,b2]和[b1,b2]合并——本质是两个线程都往同一个初始列表里加了元素,combiner又把这个已经被修改过的列表和自身合并,自然就出现了重复元素。
- 违反
reduce的契约要求reduce的accumulator和combiner必须满足一个规则:对于任意的u和t,combiner(u, accumulator(identity, t))要等于accumulator(u, t)。但你的accumulator是直接修改传入的r(也就是identity实例)并返回,这直接打破了这个契约——因为identity已经被修改,combiner合并时会把已经添加过的元素重复加入。
解决办法
针对这个场景,最推荐用collect代替reduce,因为collect就是专门为可变容器的并行收集设计的,完美匹配你的需求:
List<String> list = Stream.of("a", "b").parallel() .map(e -> List.of(e + "1", e + "2")) .collect(ArrayList::new, ArrayList::addAll, ArrayList::addAll); System.out.println(list);
这里三个参数的作用清晰明确:
ArrayList::new:每个工作线程都会创建独立的ArrayList实例,彻底避免共享修改的问题ArrayList::addAll:把每个map后的列表元素添加到当前线程的容器中ArrayList::addAll:把不同线程的容器合并为最终结果
如果一定要坚持用reduce,那必须保证accumulator和combiner都返回新的列表实例,绝不修改传入的参数:
List<String> list = Stream.of("a", "b").parallel() .map(e -> List.of(e + "1", e + "2")) .reduce(new ArrayList<>(), (r, l) -> { log.info("accumulator r:{}, l:{}", r, l); ArrayList<String> newList = new ArrayList<>(r); newList.addAll(l); return newList; }, (r1, r2) -> { log.info("combiner r:{}, l:{}", r1, r2); ArrayList<String> newList = new ArrayList<>(r1); newList.addAll(r2); return newList; }); System.out.println(list);
这样每个步骤都会创建新列表,不会共享同一个实例,也就不会出现重复元素的随机异常了。
内容的提问来源于stack exchange,提问作者waste material
相关产品推荐
相关产品推荐

