Fork/Join框架是Java7提供的一个用于并行执行任务的框架,是一个把大任务分割成若干个小任务最终汇总每个小任务结果后得到大任务结果的框架。

工作窃取算法
工作窃取算法是指某个线程从其他队列里窃取任务来执行。那么,为什么需要使用工作窃取算法呢?假如我们需要做一个比较大的任务,可以把这个任务分割为若干个互不依赖的子任务,线程们将自己的任务干完之后,与其等着不如去帮忙干活,于是就从其他任务的队列里面窃取一个任务来执行。这时他们会访问同一个队列,所以为了减少窃取任务线程和被窃取任务线程之间的竞争,通常会使用双端队列,「被窃取任务线程永远在双端队列的头部拿任务执行,而窃取任务永远在双端队列的尾部拿任务执行」。
优点:充分利用线程进行并行计算,减少了线程间的竞争 缺点:在某些情况下还是存在竞争,比如双端队列只有一个任务的时候。并且该算法会消耗了更多的系统资源,比如创建了多个线程和多个双端队列。
Fork/Join框架的设计
明白了Fork/Join框架的需求之后,我们应该如何设计一个Fork/Join框架呢?
分割任务
首先我们要有一个fork类来将大任务分割成子任务,有可能子任务还是很大,所以还需要不停地分割,直到子任务足够小
执行任务并合并结果
分割地子任务分别放在多个双端队列里,然后启动线程分别从这些双端队列里获取任务执行。子任务执行完地结果统一放在一个队列里,启动一个线程从队列里拿数据,然后合并这些数据。
Fork/Join使用两个类完成以上两件事
ForkJoinTask
首先创建一个ForkJoin任务。它提供在任务中执行fork()
和join()
操作的机制。通常情况下,我们不需要直接继承ForkJoinTask类,只需要继承它的子类:
RecursiveAction
:用于没有返回结果的任务RecursiveTask
:用于返回有返回结果的任务
ForkJoinPool
ForkJoinTask需要通过ForkJoinPool来执行
任务分割出的子任务会添加到当前工作线程所维护的双端队列中,进入队列的头部。当一个工作线程的队列里暂时没有任务时,它会随机从其他工作线程的队列的尾部获取一个任务
使用Fork/Join框架
public class CountTask extends RecursiveTask<Integer> {
private static final int THRESHOLD = 2;
private int start;
private int end;
public CountTask(int start, int end) {
this.end = end;
this.start = start;
}
@Override
protected Integer compute() {
int sum = 0;
// 如果任务足够小就计算任务
boolean canCompute = (end - start) <= THRESHOLD;
if (canCompute) {
for (int i=start; i<=end; i++) {
sum += i;
}
}
else {
int middle = ((end - start) >> 2) + start;
CountTask leftTask = new CountTask(start, middle);
CountTask rightTask = new CountTask(middle+1, end);
//执行子任务
leftTask.fork();
rightTask.fork();
//等待子任务执行完,并得到结果
int leftResult = leftTask.join();
int rightResult = rightTask.join();
//合并子任务
sum = leftResult + rightResult;
}
return sum;
}
public static void main(String[] args) {
ForkJoinPool forkJoinPool = new ForkJoinPool();
//生成一个计算任务,负责计算1+2+3+4
CountTask task = new CountTask(1, 4);
//执行一个子任务
Future<Integer> result = forkJoinPool.submit(task);
try {
System.out.println(result.get());
} catch (InterruptedException | ExecutionException e) {
e.printStackTrace();
}
}
}
Fork/Join框架的异常处理
ForkJoinTask在执行的时候可能抛出异常,但是我们没办法再主线程里直接捕获异常,所以ForkJoinTask提供了isCompletedAbnormally
方法检查任务是否已经抛出异常或已经被取消了,并且可以通过ForkJoinTask.getException
方法获取异常。
if(task.isCompletedAbnormally()) {
System.out.println(task.getException);
}
getException
方法返回Throwable
对象,如果任务被取消了则返回CancellationException
。如果任务没有完成或者没有抛出异常则返回null
Fork/Join框架的实现原理
ForkJoinPool由ForkJoinTask
数组和ForkJoinWorkerThread
数组组成,ForkJoinTask数组负责将存放程序提交给ForkJoinPool的任务,而ForkJoinWorkerThread数组则负责执行这些任务
ForkJoinTask的fork方法实现原理
public final ForkJoinTask<V> fork() {
Thread t;
if ((t = Thread.currentThread()) instanceof ForkJoinWorkerThread)
((ForkJoinWorkerThread)t).workQueue.push(this);
else
ForkJoinPool.common.externalPush(this);
return this;
}




