用处
分支/合并框架的目的是以递归方式将可以并行的任务拆分成更小的任务,然后将每个子任务的结果合并起来生成整体结果。它是ExecutorService接口的一个实现,它把子任务分配给线程池(称为ForkJoinPool)中的工作线程。
使用RecursiveTask
要把任务提交到这个池,必须创建RecursiveTask<R>的一个子类,其中R是并行化任务(以及所有子任务)产生的结果类型,或者如果任务不返回结果,则是RecursiveAction类型(当然它可能会更新其他非局部机构)。要定义RecursiveTask,只需实现它唯一的抽象方法compute:
protected abstract R compute();
这个方法同时定义了将任务拆分成子任务的逻辑,以及无法再拆分或不方便再拆分时,生成单个子任务结果的逻辑。正由于此,这个方法的实现类似于下面的伪代码:
if(任务足够小或不可分){顺序计算该任务} else {将任务分成两个子任务递归调用本方法,拆分每个子任务,等待所有子任务完成合并每个子任务的结果}

让我们试着用这个框架为一个数字范围(这里用一个long[]数组表示)求和。
import java.util.concurrent.RecursiveTask;import java.util.concurrent.ForkJoinTask;import java.util.stream.LongStream;import java.util.concurrent.ForkJoinPool;import java.util.function.*;class ForkJoinSumCalculator extends RecursiveTask<Long> {public static final long THRESHOLD = 10_000;public static final ForkJoinPool FORK_JOIN_POOL = new ForkJoinPool();private final long[] numbers;private final int start;private final int end;public ForkJoinSumCalculator(long[] numbers) {this(numbers, 0, numbers.length);}private ForkJoinSumCalculator(long[] numbers, int start, int end) {this.numbers = numbers;this.start = start;this.end = end;}@Overrideprotected Long compute() {int length = end - start;if (length <= THRESHOLD) {return computeSequentially();}ForkJoinSumCalculator leftTask = new ForkJoinSumCalculator(numbers, start, start + length/2);leftTask.fork();ForkJoinSumCalculator rightTask = new ForkJoinSumCalculator(numbers, start + length/2, end);Long rightResult = rightTask.compute();Long leftResult = leftTask.join();return leftResult + rightResult;}private long computeSequentially() {long sum = 0;for (int i = start; i < end; i++) {sum += numbers[i];}return sum;}public static long forkJoinSum(long n) {long[] numbers = LongStream.rangeClosed(1, n).toArray();ForkJoinTask<Long> task = new ForkJoinSumCalculator(numbers);return FORK_JOIN_POOL.invoke(task);}public static long measureSumPerf(Function<Long, Long> adder, long n) {long fastest=Long.MAX_VALUE;for (int i=0; i < 10; i++) {long start=System.nanoTime();long sum=adder.apply(n);long duration=(System.nanoTime() - start) 1000000;System.out.println("Result: "+sum);if (duration < fastest) fastest=duration;}return fastest;}public static void main(String[] args) {System.out.println("ForkJoin sum done in:"+measureSumPerf(ForkJoinSumCalculator::forkJoinSum, 10000000)+" msecs");}}
执行结果如下:
Result: 50000005000000Result: 50000005000000Result: 50000005000000Result: 50000005000000Result: 50000005000000Result: 50000005000000Result: 50000005000000Result: 50000005000000Result: 50000005000000Result: 50000005000000ForkJoin sum done in:92 msecs
使用分支/合并框架的最佳做法
对一个任务调用join方法会阻塞调用方,直到该任务做出结果。因此,有必要在两个子任务的计算都开始之后再调用它。否则,你得到的版本会比原始的顺序算法更慢更复杂,因为每个子任务都必须等待另一个子任务完成才能启动。
不应该在RecursiveTask内部使用ForkJoinPool的invoke方法。相反,你应该始终直接调用compute或fork方法,只有顺序代码才应该用invoke来启动并行计算。
对子任务调用fork方法可以把它排进ForkJoinPool。同时对左边和右边的子任务调用它似乎很自然,但这样做的效率要比直接对其中一个调用compute低。这样做你可以为其中一个子任务重用同一线程,从而避免在线程池中多分配一个任务造成的开销。
调试使用分支/合并框架的并行计算可能有点棘手。特别是你平常都在你喜欢的IDE里面看栈跟踪(stack trace)来找问题,但放在分支-合并计算上就不行了,因为调用compute的线程并不是概念上的调用方,后者是调用fork的那个。
和并行流一样,你不应理所当然地认为在多核处理器上使用分支/合并框架就比顺序计算快。
工作窃取
分支/合并框架工程用一种称为工作窃取(work stealing)的技术来解决让所有CPU内核都同样繁忙。在实际应用中,这意味着这些任务差不多被平均分配到ForkJoinPool中的所有线程上。每个线程都为分配给它的任务保存一个双向链式队列,每完成一个任务,就会从队列头上取出下一个任务开始执行。基于前面所述的原因,某个线程可能早早完成了分配给它的所有任务,也就是它的队列已经空了,而其他的线程还很忙。这时,这个线程并没有闲下来,而是随机选了一个别的线程,从队列的尾巴上“偷走”一个任务。这个过程一直继续下去,直到所有的任务都执行完毕,所有的队列都清空。这就是为什么要划成许多小任务而不是少数几个大任务,这有助于更好地在工作线程之间平衡负载。





