并行流选择ForkJoinPool的机制及相关技术疑问
并行流与ForkJoinPool的常见问题解析
示例代码
import java.time.Instant; import java.util.concurrent.ForkJoinPool; import java.util.concurrent.locks.LockSupport; import java.util.stream.Stream; public final class Main { public static void main(String[] args) throws Exception { log("using common FJP"); work(); log("using custom FJP"); ForkJoinPool forkJoinPool = new ForkJoinPool(4); try { forkJoinPool.submit(() -> work()).get(); } finally { forkJoinPool.shutdown(); } } private static void work() { Stream .of(0, 1, 2, 3) .parallel() .map(index -> { log("# %d started", index); LockSupport.parkNanos(500_000_000L); log("# %d stopped", index); return index; }) .reduce(0, Integer::sum); } private static void log(String format, Object... args) { System.out.format(Instant.now() + " [" + Thread.currentThread().getName() + "] " + format + "%n", args); } }
运行输出
2022-08-23T19:12:51.173532Z [main] using common FJP 2022-08-23T19:12:51.300260Z [main] # 2 started 2022-08-23T19:12:51.301861Z [ForkJoinPool.commonPool-worker-3] # 1 started 2022-08-23T19:12:51.801748Z [main] # 2 stopped 2022-08-23T19:12:51.802492Z [main] # 3 started 2022-08-23T19:12:51.802570Z [ForkJoinPool.commonPool-worker-3] # 1 stopped 2022-08-23T19:12:51.803850Z [ForkJoinPool.commonPool-worker-3] # 0 started 2022-08-23T19:12:52.303566Z [main] # 3 stopped 2022-08-23T19:12:52.304333Z [ForkJoinPool.commonPool-worker-3] # 0 stopped 2022-08-23T19:12:52.304684Z [main] using custom FJP 2022-08-23T19:12:52.307470Z [ForkJoinPool-1-worker-3] # 2 started 2022-08-23T19:12:52.308441Z [ForkJoinPool-1-worker-7] # 3 started 2022-08-23T19:12:52.309039Z [ForkJoinPool-1-worker-5] # 1 started 2022-08-23T19:12:52.309276Z [ForkJoinPool-1-worker-1] # 0 started 2022-08-23T19:12:52.808504Z [ForkJoinPool-1-worker-3] # 2 stopped 2022-08-23T19:12:52.809272Z [ForkJoinPool-1-worker-7] # 3 stopped 2022-08-23T19:12:52.809541Z [ForkJoinPool-1-worker-5] # 1 stopped 2022-08-23T19:12:52.810086Z [ForkJoinPool-1-worker-1] # 0 stopped
问题解答
1. 公共FJP场景下main线程执行任务的原因及线程身份
- 任务在main线程执行的核心原因是ForkJoinPool的工作窃取机制允许调用线程临时参与任务执行:当调用并行流的线程(这里是main)并非FJP的常驻工作线程时,公共FJP会将调用线程纳入任务执行流程,直接分配任务给它,避免线程切换开销,提升执行效率。
- main线程不是公共FJP的常驻工作线程,只是临时作为工作线程参与任务。当任务执行完成后,main线程会回到原来的执行流程,不会被FJP保留为常驻线程。
- 文档依据:
ForkJoinPool的JavaDoc明确说明,当通过invoke或类似方法提交任务时,调用线程会被当作工作线程使用,直到任务完成;并行流的底层执行逻辑正是基于ForkJoinPool.invoke实现,因此main线程会临时参与任务执行。
2. 并行流如何确定使用的FJP
- 并行流通过当前线程的ForkJoinPool关联上下文来选择使用的FJP:
- 当线程是某个FJP的工作线程时,该线程会绑定所属的FJP实例;此时在该线程中启动的并行流,会自动使用这个绑定的FJP。
- 当线程未绑定任何FJP(如main线程直接调用并行流),则默认使用公共
ForkJoinPool.commonPool()。
- 在示例中,提交到自定义FJP的
Runnable由自定义FJP的工作线程执行,此时work()中的并行流检测到当前线程属于自定义FJP,因此自动使用该实例执行并行任务。 - 文档依据:
Stream.parallel()的JavaDoc说明,并行流的执行依赖于当前线程的ForkJoinPool上下文;ForkJoinTask的文档也提到,线程在FJP中执行时会继承当前的FJP关联,后续提交的任务会优先使用该FJP。
内容的提问来源于stack exchange,提问作者Volkan Yazıcı
相关产品推荐
相关产品推荐

