暂无图片
暂无图片
暂无图片
暂无图片
暂无图片

分支/合并框架简介

程序员随想 2021-06-21
344

用处

分支/合并框架的目的是以递归方式将可以并行的任务拆分成更小的任务,然后将每个子任务的结果合并起来生成整体结果。它是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;
        }


        @Override
        protected 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: 50000005000000
          Result: 50000005000000
          Result: 50000005000000
          Result: 50000005000000
          Result: 50000005000000
          Result: 50000005000000
          Result: 50000005000000
          Result: 50000005000000
          Result: 50000005000000
          Result: 50000005000000
          ForkJoin sum done in:92 msecs

          使用分支/合并框架的最佳做法

          • 对一个任务调用join方法会阻塞调用方,直到该任务做出结果。因此,有必要在两个子任务的计算都开始之后再调用它。否则,你得到的版本会比原始的顺序算法更慢更复杂,因为每个子任务都必须等待另一个子任务完成才能启动。

          • 不应该在RecursiveTask内部使用ForkJoinPool的invoke方法。相反,你应该始终直接调用compute或fork方法,只有顺序代码才应该用invoke来启动并行计算。

          • 对子任务调用fork方法可以把它排进ForkJoinPool。同时对左边和右边的子任务调用它似乎很自然,但这样做的效率要比直接对其中一个调用compute低。这样做你可以为其中一个子任务重用同一线程,从而避免在线程池中多分配一个任务造成的开销。

          • 调试使用分支/合并框架的并行计算可能有点棘手。特别是你平常都在你喜欢的IDE里面看栈跟踪(stack trace)来找问题,但放在分支-合并计算上就不行了,因为调用compute的线程并不是概念上的调用方,后者是调用fork的那个。

          • 和并行流一样,你不应理所当然地认为在多核处理器上使用分支/合并框架就比顺序计算快。

          工作窃取

          分支/合并框架工程用一种称为工作窃取(work stealing)的技术来解决让所有CPU内核都同样繁忙。在实际应用中,这意味着这些任务差不多被平均分配到ForkJoinPool中的所有线程上。每个线程都为分配给它的任务保存一个双向链式队列,每完成一个任务,就会从队列头上取出下一个任务开始执行。基于前面所述的原因,某个线程可能早早完成了分配给它的所有任务,也就是它的队列已经空了,而其他的线程还很忙。这时,这个线程并没有闲下来,而是随机选了一个别的线程,从队列的尾巴上“偷走”一个任务。这个过程一直继续下去,直到所有的任务都执行完毕,所有的队列都清空。这就是为什么要划成许多小任务而不是少数几个大任务,这有助于更好地在工作线程之间平衡负载。


          分支/合并框架使用的工作窃取算法



          文章转载自程序员随想,如果涉嫌侵权,请发送邮件至:contact@modb.pro进行举报,并提供相关证据,一经查实,墨天轮将立刻删除相关内容。

          评论