ARTICLE DETAIL

资讯详情

深耕网站视觉设计与运营推广的一线实战洞察。

ThreadPoolExecutor 提交任务时的两条主要派发路径

ThreadPoolExecutor 提交任务时的两条主要派发路径 任务作为新Worker的firstTask被新线程直接执行。任务先进入workQueue之后由某个Worker在getTask()中取出来执行。核心源码在ThreadPoolExecutor.execute()、addWorker()、Worker.run()、runWorker()、getTask()这条链上。一、先看主干execute 怎么派发任务简化后的executepublic void execute(Runnable command) { //c中有两个信息//高 3 位线程池运行状态 runState //低 29 位工作线程数量 workerCountint c ctl.get(); // 1. 工作线程数 corePoolSize直接新建 Worker并把 command 作为 firstTask if (workerCountOf(c) corePoolSize) { if (addWorker(command, true)) return; c ctl.get(); } // 2. 核心线程已满尝试把任务放入队列 if (isRunning(c) workQueue.offer(command)) { int recheck ctl.get();//两个信息赋予了recheck// 入队后线程池可能被关闭需要移除任务并拒绝 if (!isRunning(recheck) remove(command)) reject(command); // 如果发现一个 Worker 都没有就补一个 Worker但它没有 firstTask else if (workerCountOf(recheck) 0)//当前线程池的 Worker 工作线程数量。 addWorker(null, false); } // 3. 队列满了尝试创建非核心线程并把 command 作为 firstTask else if (!addWorker(command, false)) { reject(command); } }这里有两个关键点addWorker(command, true)创建核心 Worker并把你的任务作为firstTask。addWorker(command, false)创建非核心 Worker也把任务作为firstTask。workQueue.offer(command)核心线程已满任务进入队列。addWorker(null, false)创建一个没有firstTask的 Worker它启动后只能去队列里取任务。worker和线程的关系是什么线程池中的线程都会被封装成一个Worker类对象ThreadPoolExecutor维护的其实就是一组Worker对象其中用集合workers存储这些Worker对象private final class Worker extends AbstractQueuedSynchronizer implements Runnable { final Thread thread; Runnable firstTask; volatile long completedTasks; Worker(Runnable firstTask) { setState(-1); this.firstTask firstTask; this.thread getThreadFactory().newThread(this); } public void run() { runWorker(this); } }关键点Worker 不是 Thread 的子类它是一个Runnable。Worker 内部有一个final Thread thread字段。//线程和worker是绑定的创建 Worker 时会通过线程工厂创建一个线程this.thread getThreadFactory().newThread(this);注意newThread(this)把Worker 自己作为线程要执行的任务。之后在addWorker中调用//下面有它的源码t.start();线程启动后执行的就是Worker.run()。Worker.run()调用runWorker(this)进入任务循环。所以关系是Thread 启动 ↓ 执行 Worker.run() ↓ runWorker(Worker) ↓ 先执行 firstTask然后循环 getTask() 从队列取任务一个 Worker 绑定一个 ThreadThread 执行 Worker 的 run 方法。二、知识点一为什么 firstTask 一开始就是你的任务看addWorker和Worker构造private boolean addWorker(Runnable firstTask, boolean core) { // ... CAS 增加 workerCount ... Worker w new Worker(firstTask);//创建时就带有firstTaskThread t w.thread; t.start();//执行的是worker中thread的任务// ... }Worker本身是一个Runnableprivate final class Worker extends AbstractQueuedSynchronizer implements Runnable { Runnable firstTask; Worker(Runnable firstTask) { setState(-1); this.firstTask firstTask; this.thread getThreadFactory().newThread(this); } public void run() { runWorker(this);//runworker才是真正干活的地方} }线程启动后执行Worker.run()进而执行runWorker(this)final void runWorker(Worker w) { Runnable task w.firstTask;//提取worker中的任务也就是创建worker时自带的任务w.firstTask null; while (task ! null || (task getTask()) ! null) { try { beforeExecute(wt, task); task.run(); afterExecute(task, null); } finally { task null; } } }注意runWorker的第一行Runnable task w.firstTask;所以如果新建Worker时传入了new Worker(command)那么这个新线程启动后第一轮循环直接执行command根本不需要经过workQueue。这就是你说的如果新建 Worker 时直接把你的任务作为 firstTask 传进去那么 task 一开始就是你的任务。更准确地说firstTask是这个新 Worker 的“启动任务”。runWorker会优先执行firstTask。执行完firstTask后才会进入getTask()循环去队列里取后续任务。三、知识点二任务先入队已有 Worker 怎么通过 getTask 取当核心线程已满execute会尝试workQueue.offer(command)如果入队成功你的任务就变成队列中的一个元素。此时提交线程的任务基本结束它不会直接执行任务也不会直接调用getTask()。真正调用getTask()的是Worker 线程自己。看runWorker的循环条件while (task ! null || (task getTask()) ! null) { task.run(); task null; }执行完当前任务后task被置为null下一轮就会调用task getTask();getTask()的核心逻辑private Runnable getTask() { for (;;) { int c ctl.get(); int rs runStateOf(c);//高 3 位线程池运行状态 runState // 线程池关闭且队列为空Worker 退出 if (rs SHUTDOWN (rs STOP || workQueue.isEmpty())) { decrementWorkerCount(); return null; } int wc workerCountOf(c); boolean timed allowCoreThreadTimeOut || wc corePoolSize; try { Runnable r timed ? workQueue.poll(keepAliveTime, TimeUnit.NANOSECONDS) : workQueue.take(); if (r ! null) return r; } catch (InterruptedException retry) { // 重试 } } }代码中rs可能对应以下状态RUNNING运行中可接收新任务也会处理队列任务SHUTDOWN不再接收新任务但仍会处理队列中已有任务STOP不再接收新任务不处理队列任务并中断正在执行的任务TIDYING所有任务已终止工作线程数为 0准备执行terminated()TERMINATEDterminated()执行完毕在你的getTask()代码里int c ctl.get(); int rs runStateOf(c); if (rs SHUTDOWN (rs STOP || workQueue.isEmpty())) { decrementWorkerCount(); return null; }含义是rs SHUTDOWN线程池已经不是RUNNING如果rs STOP直接让当前 Worker 退出不管队列里还有没有任务如果只是SHUTDOWN且workQueue.isEmpty()队列空了让 Worker 退出如果SHUTDOWN但队列非空继续取任务把已提交任务处理完boolean timed allowCoreThreadTimeOut || wc corePoolSize;allowCoreThreadTimeOut一个布尔值表示是否允许核心线程在空闲超过keepAliveTime后也被回收。默认是false。wcworkerCountOf(c)得到的当前工作线程总数。corePoolSize线程池的核心线程数。timed最终结果表示当前线程是否属于“可超时回收”的线程。整个代码的关键点核心线程默认使用workQueue.take()队列为空就阻塞等待。非核心线程或者允许核心线程超时使用workQueue.poll(keepAliveTime, ...)超时取不到就返回nullWorker 退出。当你的任务被offer进队列后阻塞队列内部会唤醒等待中的take()线程。某个空闲 Worker 被唤醒getTask()返回你的任务然后runWorker执行task.run()。所以第二句话的完整版本是如果任务先成功进入workQueue那么存活着的 Worker 在执行完当前任务后会通过getTask()从队列中取出它。如果此时没有 Worker线程池会通过addWorker(null, false)补一个 Worker这个 Worker 没有firstTask启动后也直接进入getTask()从队列取任务。也就是说不一定是“已有 Worker”也可能是“新建一个没有 firstTask 的 Worker”来消费队列。四、两条路径对比场景execute 中的动作任务如何被执行工作线程数 corePoolSizeaddWorker(command, true)新 Worker 的firstTask command启动后直接执行核心线程已满队列未满workQueue.offer(command)已有 Worker 执行完后getTask()取出若没有 Worker则addWorker(null, false)后由新 Worker 的getTask()取出核心线程已满队列已满线程数 maximumPoolSizeaddWorker(command, false)新非核心 Worker 的firstTask command启动后直接执行队列满且线程数 maximumPoolSizereject(command)拒绝策略处理五、举个具体例子假设corePoolSize 2 maximumPoolSize 4 workQueue new LinkedBlockingQueue(2)提交顺序提交task1workerCount 0 2调用addWorker(task1, true)。新 Worker-1 的firstTask task1启动后直接执行task1。提交task2workerCount 1 2调用addWorker(task2, true)。新 Worker-2 的firstTask task2启动后直接执行task2。提交task3核心线程已满workQueue.offer(task3)成功。task3进入队列。Worker-1 或 Worker-2 执行完手头任务后循环调用getTask()从队列取出task3执行。提交task4队列还能放workQueue.offer(task4)成功。之后某个空闲 Worker 通过getTask()取出执行。提交task5队列已满核心线程也满但workerCount 2 maximumPoolSize 4。调用addWorker(task5, false)创建非核心 Worker-3。Worker-3 的firstTask task5启动后直接执行task5。再提交task6队列满线程数也达到 maximumPoolSize触发拒绝策略。六、总结成一句话ThreadPoolExecutor中任务有两种主要进入执行态的方式作为新 Worker 的firstTask新线程启动后runWorker先执行firstTask所以任务一开始就被这个新 Worker 执行。先放入workQueue之后由存活 Worker 在runWorker循环中调用getTask()从阻塞队列里取出来执行如果没有 Worker则补一个firstTask null的 Worker让它通过getTask()取。所以这两句话并不矛盾它们分别对应execute()中的不同分支核心线程未满或队列满但还能扩容时走 firstTask核心线程已满且队列能入队时走 workQueue getTask。
返回列表