为何在该场景下Java线程池比parallelStream慢这么多?
本次测试基于Java 23,同时在Java 21和17上也复现了相同结果。
我近期通过JMH开展基准测试,对比线程池与parallelStream在简单计算场景下的性能差异。测试对象包括:
Executors.newFixedThreadPoolExecutors.newWorkStealingPool- 手动创建的
ForkJoinPool ForkJoinPool.commonPoolparallelStream()实现- 作为基准的串行循环
测试计算逻辑
@State(Scope.Benchmark) private static class StateData{ public static final List<ContainerClass> models = IntStream.range(1,1000_001) .mapToObj(x->{ double beamLength = 10.0; // meters // ContainerClass 是包含两个double类型参数的record return new ContainerClass(beamLength, x); }).toList(); } Callable<Double> getCallable(ContainerClass x){ return ()-> x.length()*x.load()/2.0; }
基准测试代码
@Benchmark @Fork(value = 1) @Warmup(iterations = 5) @Measurement(iterations = 5) public void testingParallelStream_toList(Blackhole bh){ var midMoments = StateData.models.parallelStream() .unordered() .map(x-> x.load()*x.length()/2.0).toList(); bh.consume(midMoments); } @Benchmark @Fork(value = 1) @Warmup(iterations = 5) @Measurement(iterations = 5) public void testingSequential(Blackhole bh) { List<Double> results = new ArrayList<>(); for(var x: StateData.models){ var result = x.load()*x.length()/2.0; results.add(result); } bh.consume(results); } @Benchmark @Fork(value = 1) @Warmup(iterations = 5) @Measurement(iterations = 5) public void testingExecutorService_FixedThreadPool(Blackhole bh) throws InterruptedException { var pool = Executors.newFixedThreadPool(Runtime.getRuntime().availableProcessors()); runPool(bh, pool); } @Benchmark @Fork(value = 1) @Warmup(iterations = 5) @Measurement(iterations = 5) public void testingExecutorService_WorkStealingPool(Blackhole bh) throws InterruptedException { var pool = Executors.newWorkStealingPool(Runtime.getRuntime().availableProcessors()); runPool(bh, pool); } @Benchmark @Fork(value = 1) @Warmup(iterations = 5) @Measurement(iterations = 5) public void testingManualFJPool(Blackhole bh) throws InterruptedException, ExecutionException { var pool = new ForkJoinPool(); runPool(bh, pool); } @Benchmark @Fork(value = 1) @Warmup(iterations = 5) @Measurement(iterations = 5) public void testingCommonFJPool(Blackhole bh) throws InterruptedException, ExecutionException { var pool = ForkJoinPool.commonPool(); List<Double> results = new ArrayList<>(); List<Callable<Double>> callables = new ArrayList<>(); for(var m: StateData.models){ callables.add(getCallable(m)); } var futures = pool.invokeAll(callables); boolean isFinished = pool.awaitQuiescence(60L, TimeUnit.MINUTES); if(!isFinished){ throw new IllegalArgumentException("Timeout"); } else { for(var future: futures){ results.add(future.get()); } } bh.consume(results); } private void runPool(Blackhole bh, ExecutorService pool) throws InterruptedException { List<Double> results = new ArrayList<>(); List<Callable<Double>> callables = new ArrayList<>(); for(var m: StateData.models){ callables.add(getCallable(m)); } var futures = pool.invokeAll(callables); // 等待线程执行完成 pool.shutdown(); try{ boolean isFinished = pool.awaitTermination(60L, TimeUnit.MINUTES); if(isFinished){ for(var future: futures){ results.add(future.get()); } bh.consume(results); } else { throw new IllegalArgumentException("所有线程完成前超时"); } } catch(Exception e){ throw new IllegalArgumentException(e); } }
测试结果
Benchmark Mode Cnt Score Error Units FunctionalVsImperative.PerfTests.testingCommonFJPool thrpt 5 10.859 ± 0.937 ops/s FunctionalVsImperative.PerfTests.testingExecutorService_FixedThreadPool thrpt 5 5.605 ± 0.907 ops/s FunctionalVsImperative.PerfTests.testingExecutorService_WorkStealingPool thrpt 5 10.278 ± 0.430 ops/s FunctionalVsImperative.PerfTests.testingManualFJPool thrpt 5 9.875 ± 1.709 ops/s FunctionalVsImperative.PerfTests.testingParallelStream_toList thrpt 5 74.648 ± 4.755 ops/s FunctionalVsImperative.PerfTests.testingSequential thrpt 5 46.895 ± 6.828 ops/s
我原本预期线程池会因parallelStream的任务拆分更优而更慢,但测试结果却出乎意料:线程池不仅远慢于parallelStream,甚至比串行循环还慢!我在1000、10000、100000次迭代下测试,趋势一致。
(注:小迭代量下串行版本比parallelStream更快符合预期,但本次大迭代量场景中parallelStream比串行快,说明计算量足够抵消线程开销)
测试环境为Core i7 10850H,12可用CPU(六核+超线程)。
编辑(23/04/25):根据讨论,我修改了State类,提前预创建Callable列表以消除实例化百万个Callable对象的开销:
@State(Scope.Benchmark) private static class StateData{ public static final List<ContainerClass> models = IntStream.range(1,1000_001) .mapToObj(x->{ double beamLength = 10.0; // meters return new ContainerClass(beamLength, x); }).toList(); public static final List<Callable<Double>> getCallables(){ List<Callable<Double>> callables = new ArrayList<>(); for(var m: StateData.models){ callables.add(StateData.getCallable(m)); } return callables; } public static Callable<Double> getCallable(ContainerClass x){ return ()-> x.length()*x.load()/2.0; } }
同时修改了runPool()方法:
private void runPool(Blackhole bh, ExecutorService pool) throws InterruptedException { List<Double> results = new ArrayList<>(); var futures = pool.invokeAll(StateData.getCallables()); // ... 其余代码不变 }
但遗憾的是,基准测试结果几乎没有改善,与之前的结果非常相似。
问题解答
核心原因分析
线程池重复创建销毁的开销
每个基准测试方法都会重新创建线程池,测试结束后立即销毁。线程池创建涉及线程初始化、锁与队列资源分配,销毁需要等待所有线程终止,百万级任务场景下这部分开销被大幅放大,直接拉低了整体性能。而parallelStream复用了ForkJoinPool.commonPool(),无需每次创建销毁,节省了大量资源。任务粒度与调度效率差异
parallelStream基于ForkJoinPool的分治策略,将大任务拆分为多个小任务,达到阈值后串行执行,调度开销远低于给每个元素单独创建Callable提交的模式。百万个独立Callable任务的提交、调度、结果收集过程,会带来极高的上下文切换和调度成本。结果收集的效率差异
parallelStream的toList()做了针对性优化:根据并行度预分配容器大小,利用ForkJoinPool特性在任务执行过程中直接填充结果,避免了百万次Future.get()的调用开销。而线程池实现需要逐个调用Future.get(),每个调用都可能涉及线程等待和锁竞争,进一步增加开销。FixedThreadPool的额外劣势
Executors.newFixedThreadPool使用无界队列,调度机制不如ForkJoinPool的工作窃取算法高效。在CPU密集型场景下,固定线程池的调度开销显著高于工作窃取池,这也是它性能最差的原因之一。
验证建议
若要公平对比,可修改测试代码:
- 将线程池创建移到
@Setup方法中,基准测试前初始化一次,测试后在@TearDown中销毁,避免重复创建销毁的开销。 - 对于
ForkJoinPool,使用ForkJoinTask而非Callable提交任务,对齐parallelStream的底层分治调度逻辑。
内容的提问来源于stack exchange,提问作者Andorrax

