如何加速并行列表处理?最优方案实测与疑问
问题: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
相关产品推荐
相关产品推荐

