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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.26 18:58:13