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

如何加速并行列表处理?最优方案实测与疑问

问题:IO密集型任务的并行处理优化

问题背景

需要处理一个字符串列表,每条字符串单独处理,单条处理耗时约40-50毫秒,200条字符串的总处理耗时已达约1秒,耗时过长。

尝试通过并行化提速,最初用CompletableFuture实现,后测试多种方案并编写基准测试,结果显示带FixedThreadPool的CompletableFuture是最优选项,但对该方案中设置100线程存在顾虑:部分任务会在队列中等待,而增加到500线程又明显不合理。

注:doWork方法模拟字符串处理逻辑,该方法因调用第三方服务产生阻塞,无法优化。

基准测试结果

Benchmark                                              Mode  Cnt  Score   Error  Units
TaskBenchmark.benchmarkCompletableFuture              thrpt    9  0,430 ± 0,001  ops/s
TaskBenchmark.benchmarkCompletableFutureWithExecutor  thrpt    9  3,848 ± 0,010  ops/s
TaskBenchmark.benchmarkForkJoin                       thrpt    9  0,437 ± 0,024  ops/s
TaskBenchmark.benchmarkParallel                       thrpt    9  0,423 ± 0,008  ops/s
TaskBenchmark.benchmarkCompletableFuture               avgt    9  2,326 ± 0,010   s/op
TaskBenchmark.benchmarkCompletableFutureWithExecutor   avgt    9  0,260 ± 0,001   s/op
TaskBenchmark.benchmarkForkJoin                        avgt    9  2,286 ± 0,122   s/op
TaskBenchmark.benchmarkParallel                        avgt    9  2,360 ± 0,018   s/op

疑问

  • 当前实现是否存在错误?
  • 有没有更优的并行处理方案?

实现代码

CompletableFuture实现

public class ComputableFutureExample implements AutoCloseable {
   private ExecutorService executor = Executors.newFixedThreadPool(100);

    public List<String> handle(List<String> records) {
        List<CompletableFuture<String>> futures = new ArrayList<>();
        for (String record : records) {
            CompletableFuture<String> future = CompletableFuture.supplyAsync(() -> WorkUtil.doWork(record));
            futures.add(future);
        }

        return CompletableFuture.allOf(futures.toArray(new CompletableFuture[0]))
                .thenApply(v -> futures.stream()
                        .map(CompletableFuture::join)
                        .collect(Collectors.toList()))
                .join();
    }

    public List<String> handleWithExecutor(List<String> records) {
        List<CompletableFuture<String>> futures = new ArrayList<>();
        for (String record : records) {
            CompletableFuture<String> future = CompletableFuture.supplyAsync(() -> WorkUtil.doWork(record), executor);
            futures.add(future);
        }

        return CompletableFuture.allOf(futures.toArray(new CompletableFuture[0]))
                .thenApply(v -> futures.stream()
                        .map(CompletableFuture::join)
                        .collect(Collectors.toList()))
                .join();
    }

    @Override
    public void close() {
        executor.shutdown();
    }
}

RecursiveTask实现

public class RecordTask extends RecursiveTask<List<String>> {
    private final List<String> records;

    public RecordTask(List<String> records) {
        this.records = records;
    }

    @Override
    protected List<String> compute() {
        if (records.size() > 1) {
            return ForkJoinTask.invokeAll(createSubtasks())
                    .stream()
                    .map(ForkJoinTask::join)
                    .flatMap(List::stream)
                    .toList();
        } else {
            return processing(records);
        }
    }

    private Collection<RecordTask> createSubtasks() {
        List<RecordTask> dividedTasks = new ArrayList<>();
        dividedTasks.add(new RecordTask(
                records.subList(0, records.size() / 2)));
        dividedTasks.add(new RecordTask(
                records.subList(records.size() / 2, records.size())));
        return dividedTasks;
    }

    private List<String> processing(List<String> records) {
        return records.stream().map(WorkUtil::doWork).toList();
    }

}

WorkUtil类

public class WorkUtil {
    public static String doWork(String record) {
        try {
            TimeUnit.MILLISECONDS.sleep(50);
        } catch (InterruptedException e) {
            e.printStackTrace();
        }
        return record + " completed";
    }
}

基准测试代码

@OutputTimeUnit(TimeUnit.SECONDS)
@Warmup(iterations = 3, time = 5)
@Measurement(iterations = 3, time = 5)
@BenchmarkMode({Mode.AverageTime, Mode.Throughput})
@Fork(value = 3)
public class TaskBenchmark {
    @State(Scope.Benchmark)
    public static class DataProvider {
        private List<String> records;

        @Setup(Level.Invocation)
        public void setup() {
            this.records = IntStream.range(0, 500).mapToObj(i -> UUID.randomUUID().toString() + i).toList();
        }
    }

    @Benchmark
    public void benchmarkForkJoin(DataProvider dataProvider) {
        RecordTask recordTask = new RecordTask(dataProvider.records);
        try (ForkJoinPool forkJoinPool = new ForkJoinPool()) {
            forkJoinPool.invoke(recordTask);
        }
    }

    @Benchmark
    public void benchmarkCompletableFuture(DataProvider dataProvider) {
        try (ComputableFutureExample example = new ComputableFutureExample();) {
            example.handle(dataProvider.records);
        }
    }

    @Benchmark
    public void benchmarkCompletableFutureWithExecutor(DataProvider dataProvider) {
        try (ComputableFutureExample example = new ComputableFutureExample();) {
            example.handleWithExecutor(dataProvider.records);
        }
    }

    @Benchmark
    public void benchmarkParallel(DataProvider dataProvider) {
        dataProvider.records.parallelStream().map(WorkUtil::doWork).toList();
    }
}

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.24 11:19:51