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

ForkJoinPool与ExecutorService对比:为何get()未阻塞致数据丢失?

为什么ForkJoinPool.get()没有阻塞等待所有批量插入完成?

我们需要向数据库插入约400万条记录,数据来自56个文件,每个文件约8万条记录。最初用ForkJoinPool按1000条批量插入,期望所有线程完成后再执行后续流程,但发现ForkJoinPool.get()并未阻塞,导致约5000条记录丢失。改用ExecutorService后所有记录正常保存,二者性能相近(均耗时约5分钟),下面分析ForkJoinPool.get()未阻塞的核心原因。

两种实现代码

ForkJoinPool实现

// 将80000条记录分割为1000条一组的批次
List<List<EntityObject>> entitiesPartitonList = Lists.partition(EntityObjectsList);

ForkJoinPool filesFJP= new ForkJoinPool(80);
filesFJP.submit(() -> entitiesPartitonList.parallelStream().forEach(entitesBatch -> {
            jdbcTemplate.batchUpdate(query, entitesBatch.toArray(new Map[entitesBatch.size()]));
})).get();

ExecutorService实现

List<Future<Boolean>> futures = new ArrayList<>();
ExecutorService es = Executors.newFixedThreadPool(80);

entitiesPartitonList.forEach(entitesBatch -> futures.add(es.submit(() -> {
        jdbcTemplate.batchUpdate(query, entitesBatch.toArray(new Map[entitesBatch.size()]));
        return true;
})));

futures.forEach(fut -> {
    try {
        fut.get();
    } catch (Exception e) {
        // 异常处理逻辑
    }
});

核心原因分析

1. parallelStream默认使用公共ForkJoinPool,而非自定义池

你在自定义的filesFJP中提交了一个任务,但任务内部的parallelStream()默认会使用ForkJoinPool.commonPool()(JVM全局公共池),而非你创建的filesFJP。这意味着:

  • 外层调用get()等待的只是「启动parallelStream遍历」这个触发任务,而非所有批量插入子任务
  • 触发任务在提交完parallelStream的遍历请求后就会结束,get()因此提前返回,此时公共池里的插入任务可能还在后台运行,后续流程提前推进就会导致部分记录未插入完成

2. parallelStream的异步执行未被外层任务等待

parallelStream().forEach()是异步分发任务到公共池的,外层的ForkJoinTask不会等待流中所有元素处理完毕。外层任务执行到parallelStream().forEach(...)时,只是把遍历任务提交给公共池,随后外层任务就完成了,所以get()会立刻返回,无法确保所有批量插入完成。

两种实现的本质差异

  • ExecutorService实现中,每个批量任务都被包装为Future,遍历调用fut.get()会逐个等待所有批量插入任务完成,确保所有数据都插入后才继续后续流程
  • ForkJoinPool实现中,错误嵌套了parallelStream的异步执行,且未正确等待所有子任务,导致外层任务提前结束

修复后的ForkJoinPool实现

如果要继续使用自定义ForkJoinPool,应该直接用池的任务提交方式,避免依赖parallelStream的默认池:

List<List<EntityObject>> entitiesPartitionList = Lists.partition(EntityObjectsList, 1000);

ForkJoinPool filesFJP = new ForkJoinPool(80);
filesFJP.submit(() -> {
    entitiesPartitionList.forEach(entitiesBatch -> {
        // 提交任务到自定义池并等待完成
        filesFJP.submit(() -> {
            jdbcTemplate.batchUpdate(query, entitiesBatch.toArray(new Map[entitiesBatch.size()]));
        }).join();
    });
}).get();

内容的提问来源于stack exchange,提问作者Karthik A.K

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.23 13:57:32