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

如何用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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.09 12:01:12