NOTE

6.14 14.ThreadPool

1. 是什么 Java的线程池框架,他提供了“任务提交”与“任务执行”分离开的机制 1.1. 为什么需要线程池 用来复用线程 - 第一,线程的创建和销毁开销比较大 - 第二,线程数量过多的话会导致cpu忙于上下文切换而不“干活” 1.2. 使用场景 - 单个任务执行的时间不能太长 - 任务数很多 2

Java创建于 更新于 historical

这是历史学习笔记,可能存在过时或不完整的理解。

1. 是什么

Java的线程池框架,他提供了“任务提交”与“任务执行”分离开的机制

1.1. 为什么需要线程池

用来复用线程

  • 第一,线程的创建和销毁开销比较大
  • 第二,线程数量过多的话会导致cpu忙于上下文切换而不“干活”

1.2. 使用场景

  • 单个任务执行的时间不能太长
  • 任务数很多

2. 如何使用

public class ThreadPoolExecutorTest
{
    public static void main(String[] args)
    {
        ThreadPoolExecutor executor = new ThreadPoolExecutor(Runtime.getRuntime().availableProcessors(),
                Runtime.getRuntime().availableProcessors() * 2,
                60,
                TimeUnit.SECONDS,
                new ArrayBlockingQueue<>(Runtime.getRuntime().availableProcessors()));
        for (int i = 0; i < 100; i++)
        {
            executor.execute(()->{
                System.out.println(Thread.currentThread().getName() + "正在运行");
            });
        }

    }
}

3. 原理分析

3.1. 工作原理图

网上的一张图总结的很好:

3.2. 状态流转图

可以参考ThreadPoolExecutor的注释

 * The runState provides the main lifecycle control, taking on values:
 *
 *   RUNNING:  Accept new tasks and process queued tasks
 *   SHUTDOWN: Don't accept new tasks, but process queued tasks
 *   STOP:     Don't accept new tasks, don't process queued tasks,
 *             and interrupt in-progress tasks
 *   TIDYING:  All tasks have terminated, workerCount is zero,
 *             the thread transitioning to state TIDYING
 *             will run the terminated() hook method
 *   TERMINATED: terminated() has completed
 *
 * The numerical order among these values matters, to allow
 * ordered comparisons. The runState monotonically increases over
 * time, but need not hit each state. The transitions are:
 *
 * RUNNING -> SHUTDOWN
 *    On invocation of shutdown(), perhaps implicitly in finalize()
 * (RUNNING or SHUTDOWN) -> STOP
 *    On invocation of shutdownNow()
 * SHUTDOWN -> TIDYING
 *    When both queue and pool are empty
 * STOP -> TIDYING
 *    When pool is empty
 * TIDYING -> TERMINATED
 *    When the terminated() hook method has completed
 *
 * Threads waiting in awaitTermination() will return when the
 * state reaches TERMINATED.
 *
 * Detecting the transition from SHUTDOWN to TIDYING is less
 * straightforward than you'd like because the queue may become
 * empty after non-empty and vice versa during SHUTDOWN state, but
 * we can only terminate if, after seeing that it is empty, we see
 * that workerCount is 0 (which sometimes entails a recheck -- see
 * below).

3.3. uml

3.3.1. 线程池框架

Executors工具类提供了静态方法创建线程池:参考Executors.md

3.3.2. 线程池返回结果

3.4. 属性

private final AtomicInteger ctl = new AtomicInteger(ctlOf(RUNNING, 0));
//29
private static final int COUNT_BITS = Integer.SIZE - 3;
//268435456
private static final int CAPACITY   = (1 << COUNT_BITS) - 1;

//线程池的运行时的状态保存在高3位中
// runState is stored in the high-order bits
//111 能接受新任务并处理已经添加的任务
private static final int RUNNING    = -1 << COUNT_BITS;
//000 不接受新任务,可以处理已经添加的任务
private static final int SHUTDOWN   =  0 << COUNT_BITS;
//001 不接受新任务,不处理已经添加的任务,中断正在处理的任务
private static final int STOP       =  1 << COUNT_BITS;
//010 所有的任务已经终止,ctl记录的任务数量为0
private static final int TIDYING    =  2 << COUNT_BITS;
//011 线程池彻底终止
private static final int TERMINATED =  3 << COUNT_BITS;

//最大的线程数
private static final int CAPACITY   = (1 << COUNT_BITS) - 1;

//获取线程池状态变量的高3bit,表示线程池状态
private static int runStateOf(int c)     { return c & ~CAPACITY; }
//获取线程池状态变量的低29位,表示线程池中当前线程数量
private static int workerCountOf(int c)  { return c & CAPACITY; }
private static int ctlOf(int rs, int wc) { return rs | wc; }

3.4.1. 解释

假设是8bit的机器,即Integer.SIZE为8,那么高3bit表示线程池状态,低5bit表示线程池中线程数量

公式 计算 结果 注释
COUNT_BITS Integer.SIZE - 3 8 - 3 = 5 0000 0101
CAPACITY (1 << COUNT_BITS) - 1 (1 << 5) - 1 = 31 0001 1111 线程池允许的最大数量
RUNNING -1 << COUNT_BITS -1 << 5 1110 0000 运行中
SHUTDOWN 0 << COUNT_BITS 0 << 5 0000 0000 关闭
STOP 1 << COUNT_BITS 1 << 5 0010 0000 停止
TIDYING 2 << COUNT_BITS 2 << 5 0100 0000 暂停
TERMINATED 3 << COUNT_BITS 3 << 5 0110 0000 终止
ctl ctlOf(int rs, int wc) { return rs | wc; } 1110 | 0 1110 0000 线程池状态+线程数量
workerCountOf workerCountOf(int c) { return c & CAPACITY; } 0001 1111
&
0000 0000
1110 0000 线程数量
runStateOf runStateOf(int c) { return c & ~CAPACITY; } 1110 0000
&
11100000
1110 0000 线程池状态

3.5. 构造函数

public ThreadPoolExecutor(int corePoolSize,
                          int maximumPoolSize,
                          long keepAliveTime,
                          TimeUnit unit,
                          BlockingQueue<Runnable> workQueue) {
    this(corePoolSize, maximumPoolSize, keepAliveTime, unit, workQueue,

         Executors.defaultThreadFactory(), defaultHandler);
}

 //全参数构造方法
 public ThreadPoolExecutor(int corePoolSize,
                          int maximumPoolSize,
                          long keepAliveTime,
                          TimeUnit unit,
                          BlockingQueue<Runnable> workQueue,
                          ThreadFactory threadFactory,
                          RejectedExecutionHandler handler) {
    if (corePoolSize < 0 ||
        maximumPoolSize <= 0 ||
        maximumPoolSize < corePoolSize ||
        keepAliveTime < 0)
        throw new IllegalArgumentException();
    if (workQueue == null || threadFactory == null || handler == null)
        throw new NullPointerException();
    this.acc = System.getSecurityManager() == null ?
            null :
            AccessController.getContext();
    //常驻线程池的线程数:线程数量<corePoolSize的时候,每提交一个任务,就会创建一个线程
    this.corePoolSize = corePoolSize;
    //允许的最大线程数:当提交的任务数超过核心线程数时并且工作队列已满时,会继续创建线程
    this.maximumPoolSize = maximumPoolSize;
    //保存任务的工作队列:当提交的任务数超过核心线程数时,继续提交会把任务保存到工作队列中
    //四种ArrayBlockingQueue(有界:基于数组),LinkedBlockingQueue(有界无界:基于链表),SynchronousQueue(无界:不存储元素),PriorityBlockingQueue(无界:优先级)
    this.workQueue = workQueue;
    //超过corePoolSize的线程如果没有执行任务,那么存活的时间
    this.keepAliveTime = unit.toNanos(keepAliveTime);
    //创建线程的工厂
    this.threadFactory = threadFactory;
    //拒绝策略:当提交的任务数超过核心线程数时并且工作队列已满并且超过最大线程数时,应该采取的策略
    //四种AbortPolicy(直接抛出异常),CallerRunsPolicy(调用者所在的线程执行),DiscardOldestPolicy(丢弃队首),DiscardPolicy(丢弃该任务)
    this.handler = handler;
}

//默认的拒绝策略:直接抛出异常
private static final RejectedExecutionHandler defaultHandler =
    new AbortPolicy();

//默认的线程创建工厂:Executors的defaultThreadFactory方法创建的是DefaultThreadFactory
public static ThreadFactory defaultThreadFactory() {
    return new DefaultThreadFactory();
}

无参构造方法就是初始化了核心属性:

  • corePoolSize:核心线程数量,线程池中应该常驻的线程数量
  • maximumPoolSize:线程池允许的最大线程数,非核心线程在超时之后会被清除
  • workQueue:阻塞队列,存储等待执行的任务
  • keepAliveTime:线程没有任务执行时可以保持的时间
  • unit:时间单位
  • threadFactory:线程工厂,来创建线程
  • rejectHandler:当拒绝任务提交时的策略(抛异常、用调用者所在的线程执行任务、丢弃队列中第一个任务执行当前任务、直接丢弃任务)

其中默认的线程创建工厂和默认的拒绝策略由下:

3.5.1. 默认的线程创建工厂

  • DefaultThreadFactory
DefaultThreadFactory() {
	SecurityManager s = System.getSecurityManager();
	//初始化线程组名
	group = (s != null) ? s.getThreadGroup() :
	                      Thread.currentThread().getThreadGroup();
	//初始化线程名前缀
	namePrefix = "pool-" +
	              poolNumber.getAndIncrement() +
	             "-thread-";
}

public Thread newThread(Runnable r) {
	Thread t = new Thread(group, r,
	                      namePrefix + threadNumber.getAndIncrement(),
	                      0);
  	//非守护
	if (t.isDaemon())
	    t.setDaemon(false);
	//优先级为NORMAL
	if (t.getPriority() != Thread.NORM_PRIORITY)
	    t.setPriority(Thread.NORM_PRIORITY);
	return t;
}
3.5.1.1. 默认的拒绝策略
  • AbortPolicy
public static class AbortPolicy implements RejectedExecutionHandler {
    public AbortPolicy() { }

    //抛出拒绝任务异常
    public void rejectedExecution(Runnable r, ThreadPoolExecutor e) {
        throw new RejectedExecutionException("Task " + r.toString() +
                                             " rejected from " +
                                             e.toString());
    }
}

3.6. execute方法

public void execute(Runnable command) {
	//执行的任务不能为空
    if (command == null)
        throw new NullPointerException();

    int c = ctl.get();//ctl存储了线程池数量(workerCountOf)和线程池状态
    //1.线程池当前线程数<corePoolSize
    if (workerCountOf(c) < corePoolSize) {
    	//创建新的线程,执行任务。成功直接返回
        if (addWorker(command, true))
            return;
        c = ctl.get();
    }
    //2.线程池中线程数>corePoolSize,那么尝试加入队列
    //线程池的状态是running(有可能其他线程调用了shutdown方法):把任务加入阻塞队列成功
    if (isRunning(c) && workQueue.offer(command)) {
    	//加入后重新检查下线程池状态(有可能其他线程调用了shutdownnow方法)
        int recheck = ctl.get();
        //线程池状态不是running,那么从队列中删除任务且成功
        if (! isRunning(recheck) && remove(command))
        	//拒绝该任务
            reject(command);
        else if (workerCountOf(recheck) == 0)//线程池中线程数目为0需要创建线程执行任务。什么时候会出现呢?
            addWorker(null, false);
    }
    //3.线程池线程数目>corePoolSize且阻塞队列已满,那么再次新建线程并运行任务
    else if (!addWorker(command, false))
    	//4.失败则拒绝该任务
        reject(command);
}

这个方法用于执行一个Runnable任务,处理流程如下:

  • 7-13行:线程池当前线程数<corePoolSize,那么创建新的线程,执行任务
  • 14-25行:线程池中线程数>corePoolSize,那么尝试加入队列
  • 27行:线程池线程数目>corePoolSize且阻塞队列已满,那么再次新建线程并运行任务
  • 29行:27行如果失败那么拒绝该任务,即线程池线程数目=maxPoolSize了,那么拒绝该任务

3.6.1. 线程池当前线程数<corePoolSize,那么创建新的线程,执行任务

int c = ctl.get();//ctl存储了线程池数量(workerCountOf)和线程池状态
//1.线程池当前线程数<corePoolSize
if (workerCountOf(c) < corePoolSize) {
	//创建新的线程,执行任务。成功直接返回
    if (addWorker(command, true))
        return;
    c = ctl.get();
}
  • 4行:线程池中线程数量还未到达corePoolSize,那么可以继续创建线程
  • 5-6行:通过addWorker创建线程并且执行任务command,成功的话返回true,然后6行返回
  • 7行:addWorker返回的false,那么重新获取线程数量

addWorker代码如下:

  • addWorker
private boolean addWorker(Runnable firstTask, boolean core) {
	//以下一大段逻辑都是处理线程池状态
    //外层循环标签
    retry:
    for (;;) {
    	//线程池状态
        int c = ctl.get();
        int rs = runStateOf(c);

		//(线程池状态比SHUTDOWN大) 且 (状态不为SHUTDOWN或者Runnable不为空或者阻塞队列为空)
		//那么直接返回false不执行
        if (rs >= SHUTDOWN &&
            ! (rs == SHUTDOWN &&
               firstTask == null &&
               ! workQueue.isEmpty()))
            return false;

		//内层循环
        for (;;) {
        	//线程池当前线程数量
            int wc = workerCountOf(c);
            //比最大线程数CAPACITY还大,不执行直接返回false
            if (wc >= CAPACITY ||
            	//core为true与corePoolSize比较,false与maximumPoolSize比较大小。大的话返回false
                wc >= (core ? corePoolSize : maximumPoolSize))
                return false;
            //增加线程池数量变量,成功退出外层循环
            if (compareAndIncrementWorkerCount(c))
                break retry;
            //设置线程池状态变量失败,并且线程池状态已经改变了,继续外层循环
            c = ctl.get();  // Re-read ctl
            if (runStateOf(c) != rs)
                continue retry;
            //设置线程池状态变量失败,并且线程池状态没有改变,继续内层循环
            // else CAS failed due to workerCount change; retry inner loop
        }
    }

	//以下一大段逻辑是在线程池状态正常的时候执行
    boolean workerStarted = false;
    boolean workerAdded = false;
    Worker w = null;
    try {
    	//创建新的线程
        w = new Worker(firstTask);
        final Thread t = w.thread;
        if (t != null) {
        	//加锁。保证保存线程的线程的hashset操作并发安全
            final ReentrantLock mainLock = this.mainLock;
            mainLock.lock();
            try {
                int rs = runStateOf(ctl.get());
				//检查线程池状态不为SHUTDOWN
                if (rs < SHUTDOWN ||
                    (rs == SHUTDOWN && firstTask == null)) {
                    if (t.isAlive()) // precheck that t is startable
                        throw new IllegalThreadStateException();
                    //线程是存放在hashset里的。private final HashSet<Worker> workers = new HashSet<Worker>();
                    workers.add(w);
                    int s = workers.size();
                    if (s > largestPoolSize)
                        largestPoolSize = s;
                    workerAdded = true;
                }
            } finally {
            	//解锁
                mainLock.unlock();
            }
            //添加成功后启动线程。运行Worker的run方法
            if (workerAdded) {
                t.start();
                workerStarted = true;
            }
        }
    } finally {
    	//处理新线程启动失败的情况
        if (! workerStarted)
            addWorkerFailed(w);
    }
    return workerStarted;
}
  • 4-37行:这段逻辑主要用于处理线程池状态。放在死循环中,只有在线程池状态正常,且28行新增线程池数量成功才能退出
  • 40-80行:线程池状态正常的时候才会执行。着重讲解的就是这一段
    • 45行:创建Worker。这个Worker是个Runnable,创建Worker实例的同时还会创建线程,至于传入的firstTask则是这个线程第一个执行的任务
    • 46行:取出这个Worker对应的线程
    • 49-50行:加锁。Worker对应的线程需要保存到HashSet中复用,但是HashSet不是线程安全的,所以需要加锁。 当然还有其他一些操作也要保证线程安全
    • 59-63行:把线程加入HashSet中,并置添加worker的标记为成功
    • 70-71行:启动这个线程,他会执行Worker的run方法
    • 77-78行:处理新线程启动失败的情况
3.6.1.1. 创建线程
//这个Worker继承了AQS,实现Runnable接口
private final class Worker
        extends AbstractQueuedSynchronizer
        implements Runnable
{
    //关联的线程
    final Thread thread;
    //Worker运行的第一个方法(构造函数传进来的)
    Runnable firstTask;
    //完成的任务数
    volatile long completedTasks;


    Worker(Runnable firstTask) {
        setState(-1); // inhibit interrupts until runWorker
        this.firstTask = firstTask;
        this.thread = getThreadFactory().newThread(this);//new Thread(Runnnable)。创建一个线程,线程启动后执行的任务是当前Worker实例的run方法,对应21行
    }


    //execute()中t.start()会调用这个方法
    public void run() {
        runWorker(this);//传入Worker实例调用runWorker方法
    }

    //继承AQS后重写的锁相关的方法
    protected boolean isHeldExclusively() {
        return getState() != 0;
    }

    protected boolean tryAcquire(int unused) {
        if (compareAndSetState(0, 1)) {
            setExclusiveOwnerThread(Thread.currentThread());
            return true;
        }
        return false;
    }

    protected boolean tryRelease(int unused) {
        setExclusiveOwnerThread(null);
        setState(0);
        return true;
    }

    public void lock()        { acquire(1); }
    public boolean tryLock()  { return tryAcquire(1); }
    public void unlock()      { release(1); }
    public boolean isLocked() { return isHeldExclusively(); }

    void interruptIfStarted() {
        Thread t;
        if (getState() >= 0 && (t = thread) != null && !t.isInterrupted()) {
            try {
                t.interrupt();
            } catch (SecurityException ignore) {
            }
        }
    }
}
  • 14-18行:构造方法创建Worker实例的同时会创建线程
  • 21-24行:上面创建的线程启动后执行的就是Worker的run方法,而这个run方法又会执行runWorker方法,如下:
3.6.1.2. 启动线程后执行的逻辑
  • runWorker
final void runWorker(Worker w) {
 	//获取当前线程以及要执行的任务
    Thread wt = Thread.currentThread();
    Runnable task = w.firstTask;//最开始执行的肯定是构造方法传入的task
    w.firstTask = null;
    //释放锁,为了防止从空的任务队列中取出任务时阻塞导致其他线程无法获取锁
    w.unlock(); 
    boolean completedAbruptly = true;
    try {
    	//不停的从队列中取出任务
        while (task != null || (task = getTask()) != null) {
            //获取锁
            w.lock();
            //线程池状态检测
            if ((runStateAtLeast(ctl.get(), STOP) ||
                 (Thread.interrupted() &&
                  runStateAtLeast(ctl.get(), STOP))) &&
                !wt.isInterrupted())
                wt.interrupt();
            try {
            	//执行任务之前执行的操作
                beforeExecute(wt, task);
                Throwable thrown = null;
                try {
                	//执行真正的任务,这里调用的是run不是start
                    task.run();
                } catch (RuntimeException x) {
                    thrown = x; throw x;
                } catch (Error x) {
                    thrown = x; throw x;
                } catch (Throwable x) {
                    thrown = x; throw new Error(x);
                } finally {
                    afterExecute(task, thrown);
                }
            } finally {
                task = null;
                //完成任务数+1
                w.completedTasks++;
                //解锁
                w.unlock();
            }
        }
        completedAbruptly = false;
    } finally {
    	//发生了异常后处理
        processWorkerExit(w, completedAbruptly);
    }
}
  • 3-5行:获取当前线程及其关联的任务。后续需要执行这个任务,也说明了每个线程一旦创建第一个执行的任务就是WOrker构造方法传入的。 w.firstTask = null不管第一个任务是否执行成功都先把他置为null,防止重复执行
  • 7行:释放锁,后续11行会调用getTask取出任务执行,如果任务队列为空那么会阻塞,所以这里解锁以免其他线程获取到了任务却无法执行
  • 11-41行:核心逻辑,执行任务。
    • 11行:从任务队列中取出任务或者该线程关联的任务不为空
    • 13行:取出了任务后需要加锁后执行
    • 26行:执行任务的run方法
    • 39行:执行完毕后完成任务数+1
  • 47行:发生了异常后,需要进行处理。由27-32行可知,任务运行期间发生了异常,捕获之后还会往外抛出,最后会进入到processWorkerExit进行处理
3.6.1.2.1. 从任务队列中拿到任务
  • getTask
private Runnable getTask() {
    boolean timedOut = false; 

	//死循环
    for (;;) {
    	//获取线程池状态
        int c = ctl.get();
        int rs = runStateOf(c);

		//必要时检查队列是否为空
        // Check if queue empty only if necessary.
        if (rs >= SHUTDOWN && (rs >= STOP || workQueue.isEmpty())) {
            decrementWorkerCount();
            return null;
        }

        int wc = workerCountOf(c);

		//大于corePoolSize的线程需要超时处理
        boolean timed = allowCoreThreadTimeOut || wc > corePoolSize;
        if ((wc > maximumPoolSize || (timed && timedOut))
            && (wc > 1 || workQueue.isEmpty())) {
            if (compareAndDecrementWorkerCount(c))
                return null;
            continue;
        }

        try {
        	//从任务队列中获取任务
        	//获取不到,会阻塞直到任务队列不为空,不会消耗CPU
            Runnable r = timed ?
                //从这里可以看出>corePoolSize的线程的超时触发机制是通过BlockingQueue完成的
                //即当workQueue为空的时候,poll方法会阻塞等待keepAliveTime的时间,超过了之后就会返回继续往下执行。然后调用take方法直接返回null,接着把timedOut设为true
                workQueue.poll(keepAliveTime, TimeUnit.NANOSECONDS) :
                workQueue.take();
            if (r != null)
                return r;
            timedOut = true;
        } catch (InterruptedException retry) {
            timedOut = false;
        }
    }
}
3.6.1.2.2. 处理任务运行【run】过程中出现异常的场景
  • processWorkerExit
private void processWorkerExit(Worker w, boolean completedAbruptly) {
    if (completedAbruptly) // If abrupt, then workerCount wasn't adjusted
        decrementWorkerCount();

	//获取锁
    final ReentrantLock mainLock = this.mainLock;
    mainLock.lock();
    try {
        completedTaskCount += w.completedTasks;
        //从hashset中删除worker
        workers.remove(w);
    } finally {
        mainLock.unlock();
    }

	//线程全部移除时需要修改线程池状态
    tryTerminate();

    int c = ctl.get();
    if (runStateLessThan(c, STOP)) {
        if (!completedAbruptly) {
            int min = allowCoreThreadTimeOut ? 0 : corePoolSize;
            if (min == 0 && ! workQueue.isEmpty())
                min = 1;
            if (workerCountOf(c) >= min)
                return; // replacement not needed
        }
        addWorker(null, false);
    }
}
3.6.1.3. 处理启动【start】线程出错的场景
  • addWorkerFailed
private void addWorkerFailed(Worker w) {
    //获取锁
    final ReentrantLock mainLock = this.mainLock;
    mainLock.lock();
    try {
        if (w != null)
        	//从hashset中删除该worker
            workers.remove(w);
        decrementWorkerCount();
        tryTerminate();
    } finally {
        mainLock.unlock();
    }
}

3.6.2. 线程池中线程数>corePoolSize,那么尝试加入队列

if (isRunning(c) && workQueue.offer(command))
//...

调用BlockingQueue的offer方法,这个方法不阻塞。添加失败直接返回false

3.6.3. 线程池线程数目>corePoolSize且阻塞队列已满,那么再次新建线程并运行任务

逻辑同上面的:线程池当前线程数<corePoolSize,那么创建新的线程,执行任务

3.6.4. 线程池线程数目>=maxPoolSize了,那么拒绝该任务

//3.线程池线程数目>corePoolSize且阻塞队列已满,那么再次新建线程并运行任务
else if (!addWorker(command, false))
	//4.失败则拒绝该任务
    reject(command);

从下面的addWorker方法有一段逻辑可以看出线程池中线程数目大于等于maximumPoolSize返回false,那么执行reject(command);

if (wc >= CAPACITY ||
	//core为true与corePoolSize比较,false与maximumPoolSize比较大小。大的话返回false
    wc >= (core ? corePoolSize : maximumPoolSize))
    return false;
3.6.4.1. 怎么拒绝的
final void reject(Runnable command) {
    //执行拒绝策略的rejectedExecution方法
    handler.rejectedExecution(command, this);
}

比如上面构造方法初始化的默认的拒绝策略就是抛出异常

3.6.4.1.1. 各种拒绝策略

RejectedExecutionHandler.md

3.7. submit方法

  • AbstractExecutorService submit
public Future<?> submit(Runnable task) {
    if (task == null) throw new NullPointerException();
    //封装成FutureTask
    RunnableFuture<Void> ftask = newTaskFor(task, null);//new FutureTask<T>(runnable, value);
    //转调java.util.concurrent.ThreadPoolExecutor#execute方法执行
    execute(ftask);
    return ftask;
}

submit方法就分两步执行:

  • 把Runnable封装成FutureTask
  • 转调java.util.concurrent.ThreadPoolExecutor#execute方法执行。这个方法之前分析过了,唯一不同的在于FutureTask的执行,执行结果通过返回的FutureTask的get方法获取。

3.7.1. 提交的任务封装成FutureTask

  • AbstractExecutorService newTaskFor
protected <T> RunnableFuture<T> newTaskFor(Runnable runnable, T value) {
    return new FutureTask<T>(runnable, value);
}

3.7.2. 执行FutureTask

  • FutureTask run
public void run() {
	//状态不为NEW 或者 CAS设置Thread失败,不执行直接返回
    if (state != NEW ||
        !UNSAFE.compareAndSwapObject(this, runnerOffset,
                                     null, Thread.currentThread()))
        return;
    try {
        Callable<V> c = callable;
        if (c != null && state == NEW) {
            V result;
            boolean ran;
            try {
            	//执行Callable的call方法
                result = c.call();
                ran = true;
            } catch (Throwable ex) {
                result = null;
                ran = false;
                setException(ex);
            }
            //执行完了把结果设到outcome变量。
            //唤醒其他等待结果的线程(如调用了get方法的主线程)
            if (ran)
                set(result);
        }
    } finally {
        runner = null;
        int s = state;
        if (s >= INTERRUPTING)
            handlePossibleCancellationInterrupt(s);
    }
}

提交的任何执行完毕后,我们需要手动

3.7.3. 获取FutureTask执行结果

  • FutureTask#get()
public V get() throws InterruptedException, ExecutionException {
    int s = state;
    //没有运行完
    if (s <= COMPLETING)
        s = awaitDone(false, 0L);
    //返回结果
    return report(s);
}
3.7.3.1. 没有运行完则阻塞或自旋等待
  • awaitDone
 private int awaitDone(boolean timed, long nanos)
    throws InterruptedException {
    final long deadline = timed ? System.nanoTime() + nanos : 0L;
    WaitNode q = null;
    boolean queued = false;
	//死循环
    for (;;) {
        if (Thread.interrupted()) {
            removeWaiter(q);
            throw new InterruptedException();
        }
		//判断FutureTask的状态
        int s = state;
        //大于COMPLETING,说明已经执行完,直接返回
        if (s > COMPLETING) {
            if (q != null)
                q.thread = null;
            return s;
        }
        //执行完了让出CPU,让其他线程可以修改状态为NORMAL
        else if (s == COMPLETING) // cannot time out yet
            Thread.yield();
        //第一次循环:没有执行完,构造WAITNODE准备加入队列阻塞
        else if (q == null)
            q = new WaitNode();
        //第二次循环:将WaitNode加入阻塞队列
        else if (!queued)
            queued = UNSAFE.compareAndSwapObject(this, waitersOffset,
                                                 q.next = waiters, q);
        else if (timed) {
            nanos = deadline - System.nanoTime();
            if (nanos <= 0L) {
                removeWaiter(q);
                return state;
            }
            LockSupport.parkNanos(this, nanos);
        }
        //第三次循环:阻塞当前线程
        else
            LockSupport.park(this);
    }
}

3.8. shutdown方法

public void shutdown() {
	//加锁
    final ReentrantLock mainLock = this.mainLock;
    mainLock.lock();
    try {
    	//安全检查
        checkShutdownAccess();
    	//修改为SHUTDOWN状态
        advanceRunState(SHUTDOWN);
        // 中断空闲的线程
        interruptIdleWorkers();
        //ScheduledThreadPoolExecutor的钩子方法
        onShutdown(); // hook for ScheduledThreadPoolExecutor
    } finally {
    	//解锁
        mainLock.unlock();
    }
    //把线程池修改为TERMINATED
    tryTerminate();
}
  • 3-4行:加锁
  • 9行:修改线程池为SHUTDOWN状态
  • 11行:中断空闲的线程,也意味着正在运行任务的线程不会终止
  • 19行:把线程池修改为TERMINATED

3.8.1. 修改线程池为SHUTDOWN状态

//传入的是SHUTDOWN
 private void advanceRunState(int targetState) {
    //死循环直到成功
    for (;;) {
        //包含了线程池状态和线程数量
        int c = ctl.get();
            //线程池状态至少为targetState(这里是SHUTDOWN)
        if (runStateAtLeast(c, targetState) ||
        //CAS修改为targetState(这里是SHUTDOWN)成功
            ctl.compareAndSet(c, ctlOf(targetState, workerCountOf(c))))
            break;
    }
}

3.8.2. 中断空闲的线程【正在运行任务的线程不会终止】

  • interruptIdleWorkers
private void interruptIdleWorkers() {
    //调用另一个interruptIdleWorkers方法,
    //传入true的话那么只中断一个不继续了
    //否则中断所有线程
    interruptIdleWorkers(false);
}
  • interruptIdleWorkers
private void interruptIdleWorkers(boolean onlyOne) {
    //先加锁
    final ReentrantLock mainLock = this.mainLock;
    mainLock.lock();
    try {
        //遍历所有Worker
        for (Worker w : workers) {
            Thread t = w.thread;
            //如果没有中断 且 AQS尝试加锁成功
            if (!t.isInterrupted() && w.tryLock()) {
                try {
                    //中断这个线程
                    t.interrupt();
                } catch (SecurityException ignore) {
                } finally {
                    w.unlock();
                }
            }
            //传入true的话那么只中断一个不继续了
            if (onlyOne)
                break;
        }
    } finally {
        //解锁
        mainLock.unlock();
    }
}

3.8.3. 把线程池修改为TERMINATED

  • tryTerminate
final void tryTerminate() {
    for (;;) {
        int c = ctl.get();
        //线程池状态为RUNNING
        if (isRunning(c) ||
        //或者 状态为TIDYING
            runStateAtLeast(c, TIDYING) ||
            //或者 状态为SHUTDOWN并且任务队列不为空
            (runStateOf(c) == SHUTDOWN && ! workQueue.isEmpty()))
            //以上三种情况那就直接返回
            return;
        //线程数量不为0,那么终止其中一个线程
        if (workerCountOf(c) != 0) { // Eligible to terminate
            interruptIdleWorkers(ONLY_ONE);
            return;
        }

        //加锁
        final ReentrantLock mainLock = this.mainLock;
        mainLock.lock();
        try {
            //把线程池状态从TIDYING修改为0
            if (ctl.compareAndSet(c, ctlOf(TIDYING, 0))) {
                try {
                    terminated();
                } finally {
                    ctl.set(ctlOf(TERMINATED, 0));
                    termination.signalAll();//这个唤醒的是啥?
                }
                return;
            }
        } finally {
            mainLock.unlock();
        }
        // else retry on failed CAS
    }
}

3.9. shutdownNow方法

public List<Runnable> shutdownNow() {
    List<Runnable> tasks;
    final ReentrantLock mainLock = this.mainLock;
    mainLock.lock();
    try {
        checkShutdownAccess();
    	//修改为STOP状态
        advanceRunState(STOP);
        // 中断所有线程
        interruptWorkers();
        // 清空任务列表
        tasks = drainQueue();
    } finally {
        mainLock.unlock();
    }
    //把线程池修改为TERMINATED    
    tryTerminate();
    return tasks;
}
  • 4行:加锁
  • 8行:修改线程池为STOP状态
  • 10行:中断所有线程,正在运行任务的线程也会终止
  • 12行:返回等待执行的任务列表
  • 16行:把线程池修改为TERMINATED

3.9.1. 修改线程池为STOP状态

//传入的是STOP
 private void advanceRunState(int targetState) {
    //死循环直到成功
    for (;;) {
        //包含了线程池状态和线程数量
        int c = ctl.get();
            //线程池状态至少为targetState(这里是STOP)
        if (runStateAtLeast(c, targetState) ||
        //CAS修改为targetState(这里是STOP)成功
            ctl.compareAndSet(c, ctlOf(targetState, workerCountOf(c))))
            break;
    }
}

3.9.2. 中断所有线程【正在运行任务的线程也会终止】

  • interruptWorkers
private void interruptWorkers() {
    //加锁
    final ReentrantLock mainLock = this.mainLock;
    mainLock.lock();
    try {
        //遍历所有Worker
        for (Worker w : workers)
            //调用interruptIfStarted方法
            w.interruptIfStarted();
    } finally {
        mainLock.unlock();
    }
}
void interruptIfStarted() {
    Thread t;
    if (getState() >= 0 && (t = thread) != null && !t.isInterrupted()) {
        try {
            //中断线程
            t.interrupt();
        } catch (SecurityException ignore) {
        }
    }
}

3.9.3. 清空任务列表

private List<Runnable> drainQueue() {
    BlockingQueue<Runnable> q = workQueue;
    ArrayList<Runnable> taskList = new ArrayList<Runnable>();
    //把任务队列中的任务移动到新的list中
    q.drainTo(taskList);//调用的BlockingQueue的drainTo方法

    //上面的drainTo方法可能失败
    if (!q.isEmpty()) {
        //所以只好遍历,一个个转移并删除
        for (Runnable r : q.toArray(new Runnable[0])) {
            if (q.remove(r))
                taskList.add(r);
        }
    }
    return taskList;
}

3.9.4. 把线程池修改为TERMINATED

参考上面

4. 总结

每台计算机能开启的线程数是有限的 线程创建、销毁开销大,如果能复用的效率高很多

核心参数

  • corePoolSize:核心线程数量,线程池中应该常驻的线程数量
  • maximumPoolSize:线程池允许的最大线程数,非核心线程在超时之后会被清除
  • workQueue:阻塞队列,存储等待执行的任务
  • keepAliveTime:线程没有任务执行时可以保持的时间
  • unit:时间单位
  • threadFactory:线程工厂,来创建线程
  • rejectHandler:当拒绝任务提交时的策略(抛异常、用调用者所在的线程执行任务、丢弃队列中第一个任务执行当前任务、直接丢弃任务)

5. 参考