CompletableFuture并行执行异常:线程超额、空文件及任务无法终止
问题分析
你的并行代码存在以下核心问题:
- 线程安全漏洞:
pageCount和isEndOfFile是无同步的共享变量,多线程同时读写会引发竞态条件:- 多个线程读取同一
pageCount,重复查询同页数据,生成重复文件; pageCount被多线程同时递增,导致跳页,查询不存在的页码返回空结果,生成空文件;isEndOfFile的修改无法被主线程及时感知,导致循环永远无法终止。
- 多个线程读取同一
- 循环逻辑错误:初始
isEndOfFile=false,while(isEndOfFile)不会进入循环(推测你实际代码是while(!isEndOfFile),但线程安全问题仍会导致循环失控)。 - 任务提交无限制:主线程疯狂向线程池提交任务,远超线程池容量,任务堆积同时无法正确终止。
- 未等待任务完成:主线程直接调用
future.complete(null)退出,导致后台任务可能未执行完毕。
修复方案
以下是修正后的代码,解决上述所有问题:
import java.util.concurrent.CompletableFuture; import java.util.concurrent.ForkJoinPool; import java.util.concurrent.atomic.AtomicInteger; import java.util.ArrayList; import java.util.List; public void generateStudentFiles() { // 定义并发数上限 int maxConcurrency = 5; ForkJoinPool customPool = new ForkJoinPool(maxConcurrency); AtomicInteger pageCount = new AtomicInteger(0); List<CompletableFuture<Void>> futures = new ArrayList<>(); boolean hasMoreData = true; while (hasMoreData) { int currentPage = pageCount.getAndIncrement(); CompletableFuture<Void> future = CompletableFuture.runAsync(() -> { // 步骤1:从数据库分页获取记录 List<StudentEntity> result = repo.getStudentDataWithPagination(50000, currentPage); // 处理查询结果:仅非空时写入S3 if (!CollectionUtils.isEmpty(result)) { writeResultToS3(result); } // 判断是否还有更多数据(当前页数据不满额时,标记无更多数据) if (result.size() < 50000) { synchronized (this) { hasMoreData = false; } } }, customPool); futures.add(future); // 控制并发数:任务数达到上限时,等待至少一个任务完成再继续提交 if (futures.size() >= maxConcurrency) { CompletableFuture.anyOf(futures.toArray(new CompletableFuture[0])).join(); futures.removeIf(CompletableFuture::isDone); } } // 等待所有剩余任务完成,再关闭线程池 CompletableFuture.allOf(futures.toArray(new CompletableFuture[0])).join(); customPool.shutdown(); }
关键修复点
- 线程安全的页码管理:用
AtomicInteger的getAndIncrement()确保每个线程获取唯一页码,避免重复查询或跳页。 - 并发数控制:通过任务列表和
CompletableFuture.anyOf()限制同时运行的任务数不超过5个,避免任务堆积。 - 安全的终止判断:用
synchronized块修改hasMoreData,确保只有一个线程能标记任务结束,避免循环失控。 - 空结果过滤:仅当查询结果非空时才调用写入方法,避免生成空文件。
- 等待任务完成:最后调用
allOf().join()等待所有任务执行完毕,确保数据全部写入S3后再关闭线程池。
额外优化建议
- 给每个CSV文件添加页码后缀(如
students_page_0.csv),避免文件覆盖或重复。 - 增加异常处理:在
runAsync中捕获异常,用exceptionally()或handle()处理单个任务失败的情况,避免整个流程终止。 - 若数据库支持,可先查询总记录数计算总页数,一次性提交对应数量的任务,简化循环逻辑。
内容的提问来源于stack exchange,提问作者Arun Sai
相关产品推荐
相关产品推荐

