【发布时间】:2014-02-05 10:30:59
【问题描述】:
是否可以为 Java 8 parallel stream 指定自定义线程池?我在任何地方都找不到它。
假设我有一个服务器应用程序,我想使用并行流。但是这个应用程序很大而且是多线程的,所以我想把它分开。我不希望在另一个模块的 applicationblock 任务的一个模块中运行缓慢的任务。
如果我不能为不同的模块使用不同的线程池,这意味着我不能在大多数现实世界的情况下安全地使用并行流。
试试下面的例子。有一些 CPU 密集型任务在单独的线程中执行。 这些任务利用并行流。第一个任务被破坏了,所以每一步需要 1 秒(通过线程睡眠模拟)。问题是其他线程卡住并等待中断的任务完成。这是一个人为的例子,但想象一个 servlet 应用程序和某人向共享分叉连接池提交一个长时间运行的任务。
public class ParallelTest {
public static void main(String[] args) throws InterruptedException {
ExecutorService es = Executors.newCachedThreadPool();
es.execute(() -> runTask(1000)); //incorrect task
es.execute(() -> runTask(0));
es.execute(() -> runTask(0));
es.execute(() -> runTask(0));
es.execute(() -> runTask(0));
es.execute(() -> runTask(0));
es.shutdown();
es.awaitTermination(60, TimeUnit.SECONDS);
}
private static void runTask(int delay) {
range(1, 1_000_000).parallel().filter(ParallelTest::isPrime).peek(i -> Utils.sleep(delay)).max()
.ifPresent(max -> System.out.println(Thread.currentThread() + " " + max));
}
public static boolean isPrime(long n) {
return n > 1 && rangeClosed(2, (long) sqrt(n)).noneMatch(divisor -> n % divisor == 0);
}
}
【问题讨论】:
-
自定义线程池是什么意思?有一个通用的 ForkJoinPool,但您始终可以创建自己的 ForkJoinPool 并向其提交请求。
-
提示:Java Champion Heinz Kabutz 检查了同样的问题,但影响更严重:公共分叉连接池的线程死锁。见javaspecialists.eu/archive/Issue223.html
标签: java concurrency parallel-processing java-8 java-stream