如何用CompletableFuture实现类似ExecutorService.invokeAll的异步文件加载?
问题描述
我的任务是用ExecutorService和CompletableFuture异步加载文件,同时测量2、4、8线程以及非并行场景下的执行时间。
以下是我第一个可正常运行但未使用CompletableFuture的实现:
File folderWithJson = new File(pathToFolderWithJson); ExecutorService executorService = Executors.newFixedThreadPool(16); Set<Callable<Boolean>> callables = new HashSet<>(); for(File file: Objects.requireNonNull(folderWithJson.listFiles())) { callables.add(() -> { System.out.println(Thread.currentThread().getName()); return getFineToStat(file); }); } executorService.invokeAll(callables); executorService.shutdown();
我尝试了下面的CompletableFuture写法,但不同线程数下执行时间几乎无差异(我有31个文件,总大小850MB,理论上应该有明显差异):
File folderWithJson = new File(pathToFolderWithJson); ExecutorService executorService = Executors.newFixedThreadPool(16); for(File file: Objects.requireNonNull(folderWithJson.listFiles())) { CompletableFuture.runAsync(() -> { try { getFineToStat(file); } catch (IOException e) { throw new RuntimeException(e); } }, executorService).join(); } executorService.shutdown();
请问怎么用CompletableFuture实现类似ExecutorService.invokeAll的功能?
问题根源
你当前的CompletableFuture写法没有实现并行执行——每次循环里调用.join()会阻塞当前线程,必须等当前异步任务执行完才会处理下一个文件,相当于把异步任务变成了串行执行,所以线程数变化对执行时间没有影响。
正确实现方式
要实现类似invokeAll的批量异步执行,需要先收集所有CompletableFuture实例,最后统一等待全部任务完成。
方案1:仅执行任务不收集返回值
File folderWithJson = new File(pathToFolderWithJson); ExecutorService executorService = Executors.newFixedThreadPool(16); // 收集所有异步任务 List<CompletableFuture<Void>> futures = new ArrayList<>(); for(File file: Objects.requireNonNull(folderWithJson.listFiles())) { CompletableFuture<Void> future = CompletableFuture.runAsync(() -> { try { System.out.println(Thread.currentThread().getName()); getFineToStat(file); } catch (IOException e) { throw new RuntimeException(e); } }, executorService); futures.add(future); } // 等待所有任务完成 CompletableFuture.allOf(futures.toArray(new CompletableFuture[0])).join(); executorService.shutdown();
方案2:收集任务返回结果(对应原Callable的Boolean返回值)
如果需要收集getFineToStat的返回结果,改用supplyAsync实现:
File folderWithJson = new File(pathToFolderWithJson); ExecutorService executorService = Executors.newFixedThreadPool(16); List<CompletableFuture<Boolean>> futures = new ArrayList<>(); for(File file: Objects.requireNonNull(folderWithJson.listFiles())) { CompletableFuture<Boolean> future = CompletableFuture.supplyAsync(() -> { try { System.out.println(Thread.currentThread().getName()); return getFineToStat(file); } catch (IOException e) { throw new RuntimeException(e); } }, executorService); futures.add(future); } // 等待所有任务完成并收集结果 CompletableFuture.allOf(futures.toArray(new CompletableFuture[0])).join(); List<Boolean> results = futures.stream() .map(CompletableFuture::join) .collect(Collectors.toList()); executorService.shutdown();
不同线程数的执行时间测量
可以把逻辑封装成方法,分别传入不同线程数(包括非并行场景)来测试:
public void measureExecutionTime(int threadCount) throws InterruptedException { File folderWithJson = new File(pathToFolderWithJson); ExecutorService executorService = null; long startTime = System.currentTimeMillis(); if (threadCount <= 1) { // 非并行场景:串行执行 for(File file: Objects.requireNonNull(folderWithJson.listFiles())) { try { getFineToStat(file); } catch (IOException e) { throw new RuntimeException(e); } } } else { executorService = Executors.newFixedThreadPool(threadCount); List<CompletableFuture<Void>> futures = new ArrayList<>(); for(File file: Objects.requireNonNull(folderWithJson.listFiles())) { CompletableFuture<Void> future = CompletableFuture.runAsync(() -> { try { getFineToStat(file); } catch (IOException e) { throw new RuntimeException(e); } }, executorService); futures.add(future); } CompletableFuture.allOf(futures.toArray(new CompletableFuture[0])).join(); executorService.shutdown(); executorService.awaitTermination(1, TimeUnit.MINUTES); } long endTime = System.currentTimeMillis(); System.out.printf("线程数:%d,执行时间:%d ms%n", threadCount, endTime - startTime); } // 调用示例 measureExecutionTime(1); // 非并行 measureExecutionTime(2); measureExecutionTime(4); measureExecutionTime(8);
内容的提问来源于stack exchange,提问作者Fetix
相关产品推荐
相关产品推荐

