多线程高并发编程(8)

一.概念

  Fork/Join就是将一个大任务分解(fork)成许多个独立的小任务,然后多线程并行去处理这些小任务,每个小任务处理完得到结果再进行合并(join)得到最终的结果。

  流程:任务继承RecursiveTask,重写compute方法,使用ForkJoinPool的submit提交任务,任务在某个线程中运行,工作任务中的compute方法的代码开始对任务进行分析,如果符合条件就进行任务拆分,拆分成多个子任务,每个子任务进行数据的计算或操作,得到结果返回给上一层任务开启线程进行合并,最终通过get获取整体处理结果。【只能将任务1个切分为两个,不能切分为3个或其他数量

  • ForkJoinTask:代表fork/join里面的任务类型,一般用它的两个子类RecursiveTask(任务有返回值)和RecursiveAction(任务没有返回值),任务的处理逻辑包括任务的切分都是在重写compute方法里面进行处理。只有ForkJoinTask任务可以被拆分运行和合并运行。可查看上篇Future源码分析的类图结构】【ForkJoinTask使用了模板模式进行设计,将ForkJoinTask的执行相关代码进行隐藏,通过提供抽象类(即子类RecursiveTask、RecursiveAction)暴露用户的实际业务处理。】

    • RecursiveTask:在进行exec之后会使用一个result的变量进行接受返回的结果;

      1public abstract class RecursiveTask<V> extends ForkJoinTask<V> { 2 V result; 3 protected abstract V compute(); 4 5 public final V getRawResult() { 6 return result; 7 } 8 9 protected final void setRawResult(V value) { 10 result = value; 11 } 12 protected final boolean exec() { 13 result = compute(); 14 return true; 15 } 16 17}
    • RecursiveAction:在进行exec之后没有返回结果;

      1public abstract class RecursiveAction extends ForkJoinTask<Void> { 2 3 protected abstract void compute(); 4 5 public final Void getRawResult() { return null; } 6 7 protected final void setRawResult(Void mustBeNull) { } 8 9 protected final boolean exec() { 10 compute(); 11 return true; 12 } 13 14} 
  • ForkJoinPool:fork/join框架的管理者,最原始的任务都要交给它来处理。它负责控制整个fork/join有多少个工作线程,工作线程的创建、机会都是由它来控制。它还负责workQueue队列的创建和分配,每当创建一个工作线程,它负责分配对应的workQueue,然后它把接到的活都交给工作线程去处理。是整个fork/join的容器。

    • ForkJoinPool.WorkQueue:双端队列,负责存储接收的任务;
  • ForkJoinWorkerThread:fork/join里面真正干活的”工人“,它继承了Thread,所以本质是一个线程。它有一个ForkJoinPool.WorkQueue的队列存放着它要干的活,接活之前它要向ForkJoinPool注册(registerWorker),拿到相应的workQueue,然后就从workQueue里面拿任务出来处理。它是依附于ForkJoinPool而存活,如果ForkJoinPool销毁了,它也会跟着结束。【每一个ForkJoinWorkerThread线程都具有一个独立的任务等待队列workQueue。】

    • 当使用ForkJoinPool进行submit任务提交时,创建1个workQueue将任务放进去,然后进行fork任务切分,如果切分后的任务放的进去之前的workQueue就放进去,不行就随机选取workQueue放进去,如果还放不了就创建一个新的workQueue放进去;

    public class ForkJoinWorkerThread extends Thread { final ForkJoinPool pool; final ForkJoinPool.WorkQueue workQueue; protected ForkJoinWorkerThread(ForkJoinPool pool) { super("aForkJoinWorkerThread"); this.pool = pool; this.workQueue = pool.registerWorker(this); } }

二.用法

  以前1+2+3+...+100这样的处理可以用for循环处理,现在使用fork/join来处理:从下面结果可以看到,大任务被不断的拆分成小任务,然后添加到工作线程的队列中,每个小任务都会被工作线程从队列中取出进行运行,然后每个小任务的结果的合并也由工作线程执行,然后不断的汇总成最终结果。【task通过ForkJoinPool来执行,分割的子任务添加到当前工作线程的队列中,进入队列的头部,当一个工作线程中没有任务时,会从其他工作线程的队列尾部获取一个任务。(工作窃取:当前工作线程对应的队列中没有任务了,从其他工作线程对应的队列中取出任务进行操作,然后将操作结果返还给对应队列的线程。)】

1public class MyFrokJoinTask extends RecursiveTask<Integer> { 2 private int begin; 3 private int end; 4 5 public MyFrokJoinTask(int begin, int end) { 6 this.begin = begin; 7 this.end = end; 8 } 9 10 public static void main(String[] args) throws Exception { 11 ForkJoinPool pool = new ForkJoinPool(); 12 ForkJoinTask<Integer> result = pool.submit(new MyFrokJoinTask(1, 100));//提交任务 13 System.out.println("计算的值:"+result.get());//得到最终的结果 14 15 } 16 17 @Override 18 protected Integer compute() { 19 int sum = 0; 20 if (end - begin <= 2) { 21 for (int i = begin; i <= end; i++) { 22 sum += i; 23 System.out.println("i:"+i); 24 } 25 } else { 26 MyFrokJoinTask d1 = new MyFrokJoinTask(begin, (begin + end) / 2); 27 MyFrokJoinTask d2 = new MyFrokJoinTask((begin + end) / 2+1, end); 28 d1.fork();//任务拆分 29 d2.fork();//任务拆分 30 Integer a = d1.join();//每个任务的结果 31 Integer b = d2.join();//每个任务的结果 32 sum = a + b;//汇总任务结果 33 System.out.println("sum:" + sum + ",a:" + a + ",b:" + b); 34 } 35 System.out.println("name:"+Thread.currentThread().getName()); 36 return sum; 37 } 38} 39//=========结果============ 40i:1 41i:2 42name:ForkJoinPool-1-worker-1 43i:3 44i:4 45name:ForkJoinPool-1-worker-1 46sum:10,a:3,b:7 47name:ForkJoinPool-1-worker-1 48i:5 49i:6 50i:7 51name:ForkJoinPool-1-worker-1 52sum:28,a:10,b:18 53name:ForkJoinPool-1-worker-1 54............... 55............... 56sum:91,a:28,b:63 57sum:99,a:45,b:54 58name:ForkJoinPool-1-worker-3 59name:ForkJoinPool-1-worker-1 60i:23 61i:24 62i:25 63name:ForkJoinPool-1-worker-2 64sum:135,a:63,b:72 65name:ForkJoinPool-1-worker-2 66sum:234,a:99,b:135 67name:ForkJoinPool-1-worker-3 68sum:325,a:91,b:234 69name:ForkJoinPool-1-worker-1 70sum:1275,a:325,b:950 71name:ForkJoinPool-1-worker-1 72sum:5050,a:1275,b:3775 73name:ForkJoinPool-1-worker-1 74计算的值:5050

三.分析

  ForkJoinPool

1ForkJoinPool forkJoinPool = new ForkJoinPool(); 2//Runtime.getRuntime().availableProcessors()当前操作系统可以使用的CPU内核数量 3public ForkJoinPool() { 4 this(Math.min(MAX_CAP, Runtime.getRuntime().availableProcessors()), 5 defaultForkJoinWorkerThreadFactory, null, false); 6} 7//this调用到下面这段代码 8public ForkJoinPool(int parallelism, 9 ForkJoinWorkerThreadFactory factory, 10 UncaughtExceptionHandler handler, 11 boolean asyncMode) { 12 this(checkParallelism(parallelism), //并行度 13 checkFactory(factory), //工作线程创建工厂 14 handler, //异常处理handler 15 asyncMode ? FIFO_QUEUE : LIFO_QUEUE, //任务队列出队模式 异步:先进先出,同步:后进先出 16 "ForkJoinPool-" + nextPoolId() + "-worker-"); 17 checkPermission(); 18} 19//上面的this最终调用到下面这段代码 20private ForkJoinPool(int parallelism, 21 ForkJoinWorkerThreadFactory factory, 22 UncaughtExceptionHandler handler, 23 int mode, 24 String workerNamePrefix) { 25 this.workerNamePrefix = workerNamePrefix; 26 this.factory = factory; 27 this.ueh = handler; 28 this.config = (parallelism & SMASK) | mode; 29 long np = (long)(-parallelism); // offset ctl counts 30 this.ctl = ((np << AC_SHIFT) & AC_MASK) | ((np << TC_SHIFT) & TC_MASK); 31}
  • parallelism:可并行数量,fork/join框架将依据这个并行数量的设定,决定框架内并行执行的线程数量。并行的每一个任务都会有一个线程进行处理;

  • factory:当fork/join创建一个新的线程时,同样会用到线程创建工厂。它实现了ForkJoinWorkerThreadFactory接口,使用默认的的接口实现类DefaultForkJoinWorkerThreadFactory来实现newThread方法创建一个新的工作线程;

    1public static interface ForkJoinWorkerThreadFactory { 2 /** 3 * Returns a new worker thread operating in the given pool. 4 */ 5 public ForkJoinWorkerThread newThread(ForkJoinPool pool); 6 } 7 8 static final class DefaultForkJoinWorkerThreadFactory 9 implements ForkJoinWorkerThreadFactory { 10 public final ForkJoinWorkerThread newThread(ForkJoinPool pool) { 11 return new ForkJoinWorkerThread(pool); 12 } 13 }
  • handler:异常捕获处理器。当执行的任务出现异常,并从任务中被抛出时,就会被handler捕获;

  • asyncMode:fork/join为每一个独立的工作线程准备了对应的待执行任务队列,这个任务队列是使用数组进行组合的双向队列。即可以使用先进先出的工作模式,也可以使用后进先出的工作模式;

   Fork()和Join()

  fork/join框架中提供的fork()和join()是最重要的两个方法,它们和parallelism(”可并行任务数量“)配合工作,可以导致拆分的子任务T1.1、T1.2甚至TX在fork/join中不同的运行效果(上面1+2....+100的每次运行的子任务都是不同的)。即TX子任务或等待其他已存在的线程运行关联的子任务(sum操作),或在运行TX的线程中”递归“执行其他任务(将1-50进行拆分后的子任务递归运行),或启动一个新的线程执行子任务(运行1-50另一边拆分的任务,即50-100的子任务)。

fork()用于将新创建的子任务放入当前线程的workQueue队列中,fork/join框架将根据当前正在并发执行ForkJoinTask任务的ForkJoinWorkerThread线程状态,决定是让这个任务在队列中等待,还是创建一个新的ForkJoinWorkedThread线程运行它,又或者是唤起其他正在等待任务的ForkJoinWorkerThread线程运行它。

join()用于让当前线程阻塞,直到对应的子任务完成运行并返回执行结果。或者,如果这个子任务存在于当前线程的任务等待队列workQueue中,则取出这个子任务进行”递归“执行,其目的是尽快得到当前子任务的运行结果,然后继续执行。

提交任务:

  1. sumbit的第一次提交:ForkJoinPool.submit(ForkJoinTask<T> task) -> externalPush(task) -> externalSubmit(task)

    1. submit:

      1public <T> ForkJoinTask<T> submit(ForkJoinTask<T> task) { 2 if (task == null) 3 throw new NullPointerException(); 4 externalPush(task); 5 return task; 6 } 7 8 public <T> ForkJoinTask<T> submit(Callable<T> task) { 9 ForkJoinTask<T> job = new ForkJoinTask.AdaptedCallable<T>(task); 10 externalPush(job); 11 return job; 12 } 13 14 public <T> ForkJoinTask<T> submit(Runnable task, T result) { 15 ForkJoinTask<T> job = new ForkJoinTask.AdaptedRunnable<T>(task, result); 16 externalPush(job); 17 return job; 18 } 19 20 public ForkJoinTask<?> submit(Runnable task) { 21 if (task == null) 22 throw new NullPointerException(); 23 ForkJoinTask<?> job; 24 if (task instanceof ForkJoinTask<?>) // avoid re-wrap 25 job = (ForkJoinTask<?>) task; 26 else 27 job = new ForkJoinTask.AdaptedRunnableAction(task); 28 externalPush(job); 29 return job; 30 }
    2. externalPush:将任务添加到随机选取的队列中或新创建的队列中;

      1final void externalPush(ForkJoinTask<?> task) { 2 WorkQueue[] ws; WorkQueue q; int m; 3 int r = ThreadLocalRandom.getProbe();//当前线程的一个随机数 4 int rs = runState;//当前容器的状态 5 //如果随机选取的队列还有空位置可以存放、队列加锁锁定成功,任务就放入队列中 6 if ((ws = workQueues) != null && (m = (ws.length - 1)) >= 0 && 7 (q = ws[m & r & SQMASK]) != null && r != 0 && rs > 0 && 8 U.compareAndSwapInt(q, QLOCK, 0, 1)) { 9 ForkJoinTask<?>[] a; int am, n, s; 10 if ((a = q.array) != null && 11 (am = a.length - 1) > (n = (s = q.top) - q.base)) { 12 int j = ((am & s) << ASHIFT) + ABASE; 13 U.putOrderedObject(a, j, task);//任务加入队列中 14 U.putOrderedInt(q, QTOP, s + 1);//挪动下次任务存放的槽的位置 15 U.putIntVolatile(q, QLOCK, 0);//队列解锁 16 if (n <= 1)//当前数组元素少时,进行唤醒当前线程;或者当没有活动线程或线程数较少时,添加新的线程 17 signalWork(ws, q); 18 return; 19 } 20 U.compareAndSwapInt(q, QLOCK, 1, 0);//队列解锁 21 } 22 externalSubmit(task);//升级版的externalPush 23 } 24 25 26 volatile int runState; // lockable status锁定状态 27 // runState: SHUTDOWN为负数,其他的为2的次幂 28 private static final int RSLOCK = 1; 29 private static final int RSIGNAL = 1 << 1;//唤醒 30 private static final int STARTED = 1 << 2;//启动 31 private static final int STOP = 1 << 29;//停止 32 private static final int TERMINATED = 1 << 30;//结束 33 private static final int SHUTDOWN = 1 << 31;//关闭
    3. externalSubmit:队列添加任务失败,进行升级版操作,即创建队列数组和创建队列后,将任务放入新创建的队列中;

      1private void externalSubmit(ForkJoinTask<?> task) { 2 int r; // initialize caller's probe 3 if ((r = ThreadLocalRandom.getProbe()) == 0) { 4 ThreadLocalRandom.localInit(); 5 r = ThreadLocalRandom.getProbe(); 6 } 7 for (;;) {//自旋 8 WorkQueue[] ws; WorkQueue q; int rs, m, k; 9 boolean move = false; 10 /** 11 *ForkJoinPool执行器停止工作了,抛出异常 12 *ForkJoinPool extends AbstractExecutorService 13 *abstract class AbstractExecutorService implements ExecutorService 14 *interface ExecutorService extends Executor 15 *interface Executor执行提交的对象Runnable任务 16 */ 17 if ((rs = runState) < 0) { 18 tryTerminate(false, false); // help terminate 19 throw new RejectedExecutionException(); 20 } 21 //第一次遍历,队列数组未创建,进行创建 22 else if ((rs & STARTED) == 0 || // initialize初始化 23 ((ws = workQueues) == null || (m = ws.length - 1) < 0)) { 24 int ns = 0; 25 rs = lockRunState(); 26 try { 27 if ((rs & STARTED) == 0) { 28 U.compareAndSwapObject(this, STEALCOUNTER, null, 29 new AtomicLong()); 30 // create workQueues array with size a power of two 31 int p = config & SMASK; // ensure at least 2 slots,config是CPU核数 32 int n = (p > 1) ? p - 1 : 1; 33 n |= n >>> 1; n |= n >>> 2; n |= n >>> 4; 34 n |= n >>> 8; n |= n >>> 16; n = (n + 1) << 1; 35 workQueues = new WorkQueue[n];//创建 36 ns = STARTED; 37 } 38 } finally { 39 unlockRunState(rs, (rs & ~RSLOCK) | ns); 40 } 41 } 42 //第三次遍历,把任务放入队列中 43 else if ((q = ws[k = r & m & SQMASK]) != null) { 44 if (q.qlock == 0 && U.compareAndSwapInt(q, QLOCK, 0, 1)) { 45 ForkJoinTask<?>[] a = q.array; 46 int s = q.top; 47 boolean submitted = false; // initial submission or resizing 48 try { // locked version of push 49 if ((a != null && a.length > s + 1 - q.base) || 50 (a = q.growArray()) != null) { 51 int j = (((a.length - 1) & s) << ASHIFT) + ABASE; 52 U.putOrderedObject(a, j, task); 53 U.putOrderedInt(q, QTOP, s + 1); 54 submitted = true; 55 } 56 } finally { 57 U.compareAndSwapInt(q, QLOCK, 1, 0); 58 } 59 if (submitted) { 60 signalWork(ws, q); 61 return; 62 } 63 } 64 move = true; // move on failure 65 } 66 //第二次遍历,队列数组为空,创建队列 67 else if (((rs = runState) & RSLOCK) == 0) { // create new queue 68 q = new WorkQueue(this, null); 69 q.hint = r; 70 q.config = k | SHARED_QUEUE; 71 q.scanState = INACTIVE; 72 rs = lockRunState(); // publish index 73 if (rs > 0 && (ws = workQueues) != null && 74 k < ws.length && ws[k] == null) 75 ws[k] = q; // else terminated 76 unlockRunState(rs, rs & ~RSLOCK); 77 } 78 else 79 move = true; // move if busy 80 if (move) 81 r = ThreadLocalRandom.advanceProbe(r); 82 } 83}
  2. fork任务切分的提交:ForkJoinTask.fork() -> ForkJoinWorkerThread.workQueue.push(task)/ForkJoinPool.common.externalPush(task) -> ForkJoinPool.push(task)/externalPush(task)

    1. fork:

      1public final ForkJoinTask<V> fork() { 2 Thread t; 3 if ((t = Thread.currentThread()) instanceof ForkJoinWorkerThread)//当前线程是workerThread,任务直接放入workerThread当前的workQueue 4 ((ForkJoinWorkerThread)t).workQueue.push(this); 5 else 6 ForkJoinPool.common.externalPush(this);//将任务添加到随机选取的队列中或新创建的队列中 7 return this; 8 }
    2. push:

      public class ForkJoinPool extends AbstractExecutorService { static final class WorkQueue { final void push(ForkJoinTask<?> task) { ForkJoinTask<?>[] a; ForkJoinPool p; int b = base, s = top, n; if ((a = array) != null) { // ignore if queue removed,队列被移除忽略 int m = a.length - 1; // fenced write for task visibility U.putOrderedObject(a, ((m & s) << ASHIFT) + ABASE, task);//任务加入队列中 U.putOrderedInt(this, QTOP, s + 1);//挪动下次任务存放的槽的位置 if ((n = s - b) <= 1) {//当前数组元素少时,进行唤醒当前线程;或者当没有活动线程或线程数较少时,添加新的线程 if ((p = pool) != null) p.signalWork(p.workQueues, this); } else if (n >= m)//数组所有元素都满了进行2倍扩容 growArray(); } } final ForkJoinTask<?>[] growArray() { ForkJoinTask<?>[] oldA = array; int size = oldA != null ? oldA.length << 1 : INITIAL_QUEUE_CAPACITY;//2倍扩容或初始化 if (size > MAXIMUM_QUEUE_CAPACITY) throw new RejectedExecutionException("Queue capacity exceeded"); int oldMask, t, b; ForkJoinTask<?>[] a = array = new ForkJoinTask<?>[size]; if (oldA != null && (oldMask = oldA.length - 1) >= 0 && (t = top) - (b = base) > 0) { int mask = size - 1; do { // emulate poll from old array, push to new array遍历从旧数组中取出放到新数组中 ForkJoinTask<?> x; int oldj = ((b & oldMask) << ASHIFT) + ABASE; int j = ((b & mask) << ASHIFT) + ABASE; x = (ForkJoinTask<?>)U.getObjectVolatile(oldA, oldj);//从旧数组中取出 if (x != null && U.compareAndSwapObject(oldA, oldj, x, null))//将旧数组取出的位置的对象置为null U.putObjectVolatile(a, j, x);//放入新数组 } while (++b != t); } return a; } } }

  任务的消费

  任务的消费的执行链路是ForkJoinTask.doExec() -> RecursiveTask.exec()/RecursiveAction.exec() -> 覆盖重写的compute()

  1.  doExec:任务的执行入口

    1final int doExec() { 2 int s; boolean completed; 3 if ((s = status) >= 0) { 4 try { 5 completed = exec();//消费任务 6 } catch (Throwable rex) { 7 return setExceptionalCompletion(rex); 8 } 9 if (completed) 10 s = setCompletion(NORMAL);//任务执行完设置状态为NORMAL,并唤醒其他等待任务 11 } 12 return s; 13 } 14 protected abstract boolean exec(); 15 private int setCompletion(int completion) { 16 for (int s;;) { 17 if ((s = status) < 0) 18 return s; 19 if (U.compareAndSwapInt(this, STATUS, s, s | completion)) {//任务状态修改为NORMAL 20 if ((s >>> 16) != 0)//状态不是SMASK 21 synchronized (this) { notifyAll(); }//唤醒其他等待任务 22 return completion; 23 } 24 } 25 } 26 /** The run status of this task 任务的运行状态*/ 27 volatile int status; // accessed directly by pool and workers由ForkJoinPool池或ForkJoinWorkerThread控制 28 static final int DONE_MASK = 0xf0000000; // mask out non-completion bits 29 static final int NORMAL = 0xf0000000; // must be negative 30 static final int CANCELLED = 0xc0000000; // must be < NORMAL 31 static final int EXCEPTIONAL = 0x80000000; // must be < CANCELLED 32 static final int SIGNAL = 0x00010000; // must be >= 1 << 16 33 static final int SMASK = 0x0000ffff; // short bits for tags

  任务真正执行处理逻辑

  任务提交到ForkJoinPool,最终真正的是由继承Thread的ForkJoinWorkerThread的run方法来执行消费任务的,ForkJoinWorkerThread处理哪个任务是由join来出队的;

  1. ForkJoinTask.join()

    1public final V join() { 2 int s; 3 if ((s = doJoin() & DONE_MASK) != NORMAL) 4 reportException(s); 5 return getRawResult();//得到返回结果 6 } 7 private int doJoin() { 8 int s; Thread t; ForkJoinWorkerThread wt; ForkJoinPool.WorkQueue w; 9 /** 10 * (s = status) < 0 判断任务是否已经完成,完成直接返回s 11 * 任务未完成: 12 * 1)线程是ForkJoinWorkerThread,tryUnpush任务出队然后消费任务doExec 13 * 1.1)出队或消费失败,执行awaitJoin进行自旋,如果任务状态是完成就退出,否则继续尝试出队,直到任务完成或超时为止; 14 * 2)如果线程不是ForkJoinWorkerThread,执行externalAwaitDone进行出队消费 15 */ 16 return (s = status) < 0 ? s : 17 ((t = Thread.currentThread()) instanceof ForkJoinWorkerThread) ? 18 (w = (wt = (ForkJoinWorkerThread)t).workQueue). 19 tryUnpush(this) && (s = doExec()) < 0 ? s : 20 wt.pool.awaitJoin(w, this, 0L) : 21 externalAwaitDone(); 22 } 23 private void reportException(int s) { 24 if (s == CANCELLED)//取消 25 throw new CancellationException(); 26 if (s == EXCEPTIONAL)//异常 27 rethrow(getThrowableException()); 28 }
  2. awaitJoin:

    1public class ForkJoinPool{ 2 final int awaitJoin(WorkQueue w, ForkJoinTask<?> task, long deadline) { 3 int s = 0; 4 if (task != null && w != null) { 5 ForkJoinTask<?> prevJoin = w.currentJoin; 6 U.putOrderedObject(w, QCURRENTJOIN, task); 7 CountedCompleter<?> cc = (task instanceof CountedCompleter) ? 8 (CountedCompleter<?>)task : null; 9 for (;;) { 10 if ((s = task.status) < 0)//任务完成退出 11 break; 12 if (cc != null)//当前任务即将完成,检查是否还有其他的等待任务,如果有 13 //运行当前队列的其他任务,若当前的队列中没有任务了,则窃取其他队列的任务并运行 14 helpComplete(w, cc, 0); 15 //当前队列没有任务了,或队列只剩下最后一个任务执行完了 16 else if (w.base == w.top || w.tryRemoveAndExec(task)) 17 helpStealer(w, task);//窃取其他队列的任务 18 if ((s = task.status) < 0) 19 break; 20 long ms, ns; 21 if (deadline == 0L) 22 ms = 0L; 23 else if ((ns = deadline - System.nanoTime()) <= 0L)//超时退出 24 break; 25 else if ((ms = TimeUnit.NANOSECONDS.toMillis(ns)) <= 0L) 26 ms = 1L; 27 if (tryCompensate(w)) {//当前队列阻塞了 28 task.internalWait(ms);//进行等待 29 U.getAndAddLong(this, CTL, AC_UNIT); 30 } 31 } 32 U.putOrderedObject(w, QCURRENTJOIN, prevJoin); 33 } 34 return s; 35 } 36 }
  3. externalAwaitDone:

    1private int externalAwaitDone() { 2 /** 3 * 当前任务是CountedCompleter 4 * 1)是则执行ForkJoinPool.common.externalHelpComplete() 5 * 2)否则执行ForkJoinPool.common.tryExternalUnpush(this)进行任务出队 6 * 2.1)出队成功,进行doExec()消费,否则进行阻塞等待 7 */ 8 int s = ((this instanceof CountedCompleter) ? // try helping 9 ForkJoinPool.common.externalHelpComplete( 10 (CountedCompleter<?>)this, 0) : 11 ForkJoinPool.common.tryExternalUnpush(this) ? doExec() : 0); 12 if (s >= 0 && (s = status) >= 0) {//任务未完成 13 boolean interrupted = false; 14 do { 15 if (U.compareAndSwapInt(this, STATUS, s, s | SIGNAL)) {//任务状态标记为SIGNAL 16 synchronized (this) { 17 if (status >= 0) { 18 try { 19 wait(0L);//阻塞等待 20 } catch (InterruptedException ie) {//有中断异常 21 interrupted = true;//设置中断标识为true 22 } 23 } 24 else 25 notifyAll();//任务完成唤醒其他任务 26 } 27 } 28 } while ((s = status) >= 0); 29 if (interrupted) 30 Thread.currentThread().interrupt();//当前线程进行中断 31 } 32 return s; 33 } 34 final int externalHelpComplete(CountedCompleter<?> task, int maxTasks) { 35 WorkQueue[] ws; int n; 36 int r = ThreadLocalRandom.getProbe(); 37 //没有任务直接结束,有任务则执行helpComplete 38 //helpComplete:运行随机选取的队列的任务,若选取的队列中没有任务了,则窃取其他队列的任务并运行 39 return ((ws = workQueues) == null || (n = ws.length) == 0) ? 0 : 40 helpComplete(ws[(n - 1) & r & SQMASK], task, maxTasks); 41 } 
  4. run和工作窃取

任务是由workThread来窃取的,workThread是一个线程。线程的所有逻辑都是由run()方法执行:

1public class ForkJoinWorkerThread extends Thread { 2 public void run() { 3 if (workQueue.array == null) { // only run once 4 Throwable exception = null; 5 try { 6 onStart();//初始化状态 7 pool.runWorker(workQueue);//处理任务队列 8 } catch (Throwable ex) { 9 exception = ex; 10 } finally { 11 try { 12 onTermination(exception); 13 } catch (Throwable ex) { 14 if (exception == null) 15 exception = ex; 16 } finally { 17 pool.deregisterWorker(this, exception); 18 } 19 } 20 } 21 } 22} 23 public class ForkJoinPool{ 24 final void runWorker(WorkQueue w) { 25 w.growArray(); // allocate queue,队列初始化 26 int seed = w.hint; // initially holds randomization hint 27 int r = (seed == 0) ? 1 : seed; // avoid 0 for xorShift 28 for (ForkJoinTask<?> t;;) {//自旋 29 if ((t = scan(w, r)) != null)//从队列中窃取任务成功,scan()进行任务窃取 30 w.runTask(t);//执行任务,内部方法调用了doExec()进行任务的消费 31 else if (!awaitWork(w, r))//队列没有任务了则结束 32 break; 33 r ^= r << 13; r ^= r >>> 17; r ^= r << 5; // xorshift 34 } 35 } 36 }
    1. scan:

      1private ForkJoinTask<?> scan(WorkQueue w, int r) { 2 WorkQueue[] ws; int m; 3 if ((ws = workQueues) != null && (m = ws.length - 1) > 0 && w != null) { 4 int ss = w.scanState; // initially non-negative 5 for (int origin = r & m, k = origin, oldSum = 0, checkSum = 0;;) { 6 WorkQueue q; ForkJoinTask<?>[] a; ForkJoinTask<?> t; 7 int b, n; long c; 8 if ((q = ws[k]) != null) { //随机选中了非空队列 q 9 if ((n = (b = q.base) - q.top) < 0 && 10 (a = q.array) != null) { // non-empty 11 long i = (((a.length - 1) & b) << ASHIFT) + ABASE; //从尾部出队,b是尾部下标 12 if ((t = ((ForkJoinTask<?>) 13 U.getObjectVolatile(a, i))) != null && 14 q.base == b) { 15 if (ss >= 0) { 16 if (U.compareAndSwapObject(a, i, t, null)) { //利用cas出队 17 q.base = b + 1; 18 if (n < -1) // signal others 19 signalWork(ws, q); 20 return t; //出队成功,成功窃取一个任务! 21 } 22 } 23 else if (oldSum == 0 && // try to activate 队列没有激活,尝试激活 24 w.scanState < 0) 25 tryRelease(c = ctl, ws[m & (int)c], AC_UNIT); 26 } 27 if (ss < 0) // refresh 28 ss = w.scanState; 29 r ^= r << 1; r ^= r >>> 3; r ^= r << 10; 30 origin = k = r & m; // move and rescan 31 oldSum = checkSum = 0; 32 continue; 33 } 34 checkSum += b; 35 }<br data-filtered="filtered"> //k = k + 1表示取下一个队列 如果(k + 1) & m == origin表示 已经遍历完所有队列了 36 if ((k = (k + 1) & m) == origin) { // continue until stable 37 if ((ss >= 0 || (ss == (ss = w.scanState))) && 38 oldSum == (oldSum = checkSum)) { 39 if (ss < 0 || w.qlock < 0) // already inactive 40 break; 41 int ns = ss | INACTIVE; // try to inactivate 42 long nc = ((SP_MASK & ns) | 43 (UC_MASK & ((c = ctl) - AC_UNIT))); 44 w.stackPred = (int)c; // hold prev stack top 45 U.putInt(w, QSCANSTATE, ns); 46 if (U.compareAndSwapLong(this, CTL, c, nc)) 47 ss = ns; 48 else 49 w.scanState = ss; // back out 50 } 51 checkSum = 0; 52 } 53 } 54 } 55 return null; 56 }
    2. ForkJoinPool.runTask:

      1final void runTask(ForkJoinTask<?> task) { 2 if (task != null) { 3 scanState &= ~SCANNING; // mark as busy 4 (currentSteal = task).doExec(); 5 U.putOrderedObject(this, QCURRENTSTEAL, null); // release for GC 6 execLocalTasks(); 7 ForkJoinWorkerThread thread = owner; 8 if (++nsteals < 0) // collect on overflow 9 transferStealCount(pool); 10 scanState |= SCANNING; 11 if (thread != null) 12 thread.afterTopLevelExec(); 13 } 14 }

四.总结

  对于fork/join来说,在使用时还是存在下面的一些问题的:

  • 在使用JVM的时候我们要考虑OOM的问题,如果我们的任务处理时间非常耗时,并且处理的数据非常大的时候,会造成OOM;
  • ForkJoin是通过多线程的方式进行处理任务,那么我们不得不考虑是否应该使用ForkJoin。因为当数据量不是特别大的时候,我们没有必要使用ForkJoin。因为多线程会涉及到上下文的切换,所以数据量不大的时候使用串行比使用多线程快;
    • 项目中进行本地测试发现,业务层Service进行excel表数据(数据量几百)的复杂处理,进行单线程for循环统计消耗时间,然后与使用fork/join进行处理统计消耗时间,发现fork/join的消耗时间是单线程for的2倍;
点赞
收藏

评论区

加载中...

相关推荐

MySQL:[Err] 1292 - Incorrect datetime value: ‘0000-00-00 00:00:00‘ for column ‘CREATE_TIME‘ at row 1

文章目录问题用navicat导入数据时,报错:原因这是因为当前的MySQL不支持datetime为0的情况。解决修改sql\mode:sql\mode:SQLMode定义了MySQL应支持的SQL语法、数据校验等,这样可以更容易地在不同的环境中使用MySQL。全局s

Oracle 分组与拼接字符串同时使用

SELECTT.,ROWNUMIDFROM(SELECTT.EMPLID,T.NAME,T.BU,T.REALDEPART,T.FORMATDATE,SUM(T.S0)S0,MAX(UPDATETIME)CREATETIME,LISTAGG(TOCHAR(

MySQL部分从库上面因为大量的临时表tmp_table造成慢查询

背景描述Time:20190124T00:08:14.70572408:00User@Host:@Id:Schema:sentrymetaLast_errno:0Killed:0Query_time:0.315758Lock_

皕杰报表之UUID

​在我们用皕杰报表工具设计填报报表时,如何在新增行里自动增加id呢?能新增整数排序id吗?目前可以在新增行里自动增加id,但只能用uuid函数增加UUID编码,不能新增整数排序id。uuid函数说明:获取一个UUID,可以在填报表中用来创建数据ID语法:uuid()或uuid(sep)参数说明:sep布尔值,生成的uuid中是否包含分隔符'',缺省为

手写Java HashMap源码

HashMap的使用教程HashMap的使用教程HashMap的使用教程HashMap的使用教程HashMap的使用教程22

2020年前端实用代码段,为你的工作保驾护航

有空的时候,自己总结了几个代码段,在开发中也经常使用,谢谢。1、使用解构获取json数据let jsonData  id: 1,status: "OK",data: 'a', 'b';let  id, status, data: number   jsonData;console.log(id, status, number )