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
相关产品推荐
相关产品推荐

