You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

为何在该场景下Java线程池比parallelStream慢这么多?

线程池与parallelStream的性能差异基准测试疑问

本次测试基于Java 23,同时在Java 21和17上也复现了相同结果。

我近期通过JMH开展基准测试,对比线程池与parallelStream在简单计算场景下的性能差异。测试对象包括:

  • Executors.newFixedThreadPool
  • Executors.newWorkStealingPool
  • 手动创建的ForkJoinPool
  • ForkJoinPool.commonPool
  • parallelStream()实现
  • 作为基准的串行循环

测试计算逻辑

@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());
    // ... 其余代码不变
}

但遗憾的是,基准测试结果几乎没有改善,与之前的结果非常相似。


问题解答

核心原因分析

  1. 线程池重复创建销毁的开销
    每个基准测试方法都会重新创建线程池,测试结束后立即销毁。线程池创建涉及线程初始化、锁与队列资源分配,销毁需要等待所有线程终止,百万级任务场景下这部分开销被大幅放大,直接拉低了整体性能。而parallelStream复用了ForkJoinPool.commonPool(),无需每次创建销毁,节省了大量资源。

  2. 任务粒度与调度效率差异
    parallelStream基于ForkJoinPool的分治策略,将大任务拆分为多个小任务,达到阈值后串行执行,调度开销远低于给每个元素单独创建Callable提交的模式。百万个独立Callable任务的提交、调度、结果收集过程,会带来极高的上下文切换和调度成本。

  3. 结果收集的效率差异
    parallelStream的toList()做了针对性优化:根据并行度预分配容器大小,利用ForkJoinPool特性在任务执行过程中直接填充结果,避免了百万次Future.get()的调用开销。而线程池实现需要逐个调用Future.get(),每个调用都可能涉及线程等待和锁竞争,进一步增加开销。

  4. FixedThreadPool的额外劣势
    Executors.newFixedThreadPool使用无界队列,调度机制不如ForkJoinPool的工作窃取算法高效。在CPU密集型场景下,固定线程池的调度开销显著高于工作窃取池,这也是它性能最差的原因之一。

验证建议

若要公平对比,可修改测试代码:

  • 将线程池创建移到@Setup方法中,基准测试前初始化一次,测试后在@TearDown中销毁,避免重复创建销毁的开销。
  • 对于ForkJoinPool,使用ForkJoinTask而非Callable提交任务,对齐parallelStream的底层分治调度逻辑。

内容的提问来源于stack exchange,提问作者Andorrax

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.13 09:08:10