【发布时间】:2014-10-26 18:11:12
【问题描述】:
如果池中的所有线程都在等待同一个池中的队列任务完成,则在普通线程池中会发生线程饥饿死锁。 ForkJoinPool 通过从 join() 调用内部窃取其他线程的工作来避免这个问题,而不是简单地等待。例如:
private static class ForkableTask extends RecursiveTask<Integer> {
private final CyclicBarrier barrier;
ForkableTask(CyclicBarrier barrier) {
this.barrier = barrier;
}
@Override
protected Integer compute() {
try {
barrier.await();
return 1;
} catch (InterruptedException | BrokenBarrierException e) {
throw new RuntimeException(e);
}
}
}
@Test
public void testForkJoinPool() throws Exception {
final int parallelism = 4;
final ForkJoinPool pool = new ForkJoinPool(parallelism);
final CyclicBarrier barrier = new CyclicBarrier(parallelism);
final List<ForkableTask> forkableTasks = new ArrayList<>(parallelism);
for (int i = 0; i < parallelism; ++i) {
forkableTasks.add(new ForkableTask(barrier));
}
int result = pool.invoke(new RecursiveTask<Integer>() {
@Override
protected Integer compute() {
for (ForkableTask task : forkableTasks) {
task.fork();
}
int result = 0;
for (ForkableTask task : forkableTasks) {
result += task.join();
}
return result;
}
});
assertThat(result, equalTo(parallelism));
}
但是当使用ExecutorService 接口到ForkJoinPool 时,工作窃取似乎不会发生。例如:
private static class CallableTask implements Callable<Integer> {
private final CyclicBarrier barrier;
CallableTask(CyclicBarrier barrier) {
this.barrier = barrier;
}
@Override
public Integer call() throws Exception {
barrier.await();
return 1;
}
}
@Test
public void testWorkStealing() throws Exception {
final int parallelism = 4;
final ExecutorService pool = new ForkJoinPool(parallelism);
final CyclicBarrier barrier = new CyclicBarrier(parallelism);
final List<CallableTask> callableTasks = Collections.nCopies(parallelism, new CallableTask(barrier));
int result = pool.submit(new Callable<Integer>() {
@Override
public Integer call() throws Exception {
int result = 0;
// Deadlock in invokeAll(), rather than stealing work
for (Future<Integer> future : pool.invokeAll(callableTasks)) {
result += future.get();
}
return result;
}
}).get();
assertThat(result, equalTo(parallelism));
}
粗略看一下ForkJoinPool的实现,所有常规的ExecutorService API都是使用ForkJoinTasks实现的,所以我不确定为什么会发生死锁。
【问题讨论】:
-
我不认为偷工作能避免死锁。一旦陷入僵局,就无法取得进展。工作窃取只是通过允许线程在队列为空时从其他队列窃取来避免不平衡的队列。
-
@markspace 在
ForkJoinTask的实现中,join()尝试从双端队列运行其他作业而不是停止,这样可以避免死锁。由于ForkJoinPool.invokeAll()将Callables 转换为ForkJoinTasks,我希望它也能工作。
标签: java multithreading concurrency java.util.concurrent fork-join