如何控制parallelStream的并行度?实现指定分组并行执行需求
控制ParallelStream并行度实现指定计算逻辑
要实现你想要的效果,核心是通过自定义ForkJoinPool来限制并行流的并行度,让任务按预期分组执行。具体实现如下:
实现思路
默认并行流使用JVM全局的ForkJoinPool,其并行度通常等于CPU核心数。我们可以创建并行度为2的自定义ForkJoinPool,在这个池内执行并行流操作,流会自动将元素分成两组([1,2]和[3])分别计算,最终仅触发一次combiner合并操作。
代码实现
import java.util.Arrays; import java.util.concurrent.ForkJoinPool; public class ParallelStreamDemo { public static void main(String[] args) throws Exception { // 创建并行度为2的自定义ForkJoinPool try (ForkJoinPool pool = new ForkJoinPool(2)) { Integer reducedParallel = pool.submit(() -> Arrays.asList(1, 2, 3).parallelStream() .reduce(10, (a, b) -> a + b, (a, b) -> { System.out.println("i am called"); return a + b; }) ).get(); System.out.println("结果:" + reducedParallel); } } }
执行逻辑说明
- 自定义ForkJoinPool的并行度设为2,限制了同时运行的任务数量;
- 并行流会将元素分为两组执行:
- 第一组处理元素1和2:用初始值10累加,得到
10 + 1 + 2 = 13 - 第二组处理元素3:用初始值10累加,得到
10 + 3 = 13
- 第一组处理元素1和2:用初始值10累加,得到
- 最后combiner仅执行一次,将两组结果合并为
13 + 13 = 26,同时打印一次"i am called"。
注意:使用自定义ForkJoinPool时,通过
try-with-resources语法可以确保池资源被正确释放,避免资源泄漏。
内容的提问来源于stack exchange,提问作者Tom
相关产品推荐
相关产品推荐

