😀 这里写文章的前言: 一个简单的开头,简述这篇文章讨论的问题、目标、人物、背景是什么?并简述你给出的答案。
可以说说你的故事:阻碍、努力、结果成果,意外与转折。
📝 线程池 使用线程池的好处 核心问题就是资源管理问题
频繁申请/销毁资源和调度资源,将带来额外的消耗,可能会非常巨大
对资源无限申请缺少抑制手段,容易引发系统资源耗尽的风险
系统无法合理管理内部的资源分布,会降低系统的稳定性
注意: 在 hotspot模型下,java的线程会一对一的映射内核线程,意味着每次申请和销毁都要转换到内核去操作,内核转换这个操作是十分消耗性能的,可能这个线程的时间还没 消耗+申请 加起来久
线程池的优势
提供资源的利用性: 通过池化可以重复利用已创建的线程,空闲线程可以处理新提交的任务,从而降低了创建和销毁的资源开销
提高线程的管理性: 在一个线程中管理执行任务的线程,对线程可以进行统一的创建,销毁以及监控等,对线程数做控制,防止线程无限制创建,避免线程数量的急剧上升而导致CPU过度等问题
提高程序的响应性:提交任务后,有空闲线程可以直接去执行任务,无需新建
提高系统的可扩展性:利用线程池可以更好的扩展一些功能,比如定时线程池可以实现系统的定时任务
懒惰性:先创建线程池的时候不会有任何线程,要先有第一个任务进来,才会创建线程
线程池状态 https://camo.githubusercontent.com/b79418cd3aef83df3280a4b4c450f64326ac8ffe07aa1617f602fece6f5cf713/68747470733a2f2f63646e2e6e6c61726b2e636f6d2f79757175652f302f323032332f706e672f3331363533332f313637383830393838333639322d32333765366363392d366665652d343833392d623262372d3130643362623636623536372e706e6723617665726167654875653d25323365646564656426636c69656e7449643d7531653063656562392d346330342d342666726f6d3d7061737465266865696768743d3330342669643d756666363032383961266e616d653d696d6167652e706e67266f726967696e4865696768743d333034266f726967696e57696474683d31313134266f726967696e616c547970653d62696e61727926726174696f3d3126726f746174696f6e3d302673686f775469746c653d66616c73652673697a653d3733323335267374617475733d646f6e65267374796c653d6e6f6e65267461736b49643d7535366565623462662d313561362d346564612d396361342d3138663135363265376435267469746c653d2677696474683d31313134
ctl: 对线程池的运行状态 和 线程池中有效线程的数量进行控制的一个字段
高3位保存runState,低29位保存workerCount,两个变量之间互不干扰。用一个变量去存储两个值,可避免在做相关决策时,出现不一致的情况,不必为了维护两者的一致,而占用锁资源 。通过阅读线程池源代码也可以发现,经常出现要同时判断线程池运行状态和线程数量的情况。线程池也提供了若干方法去供用户获得线程池当前的运行状态、线程个数。这里都使用的是位运算的方式,相比于基本运算,速度也会快很多 利用低29位表示线程池中线程数,高3位表示线程池的运行状态
private static int runStateOf (int c) { return c & ~CAPACITY; } private static int workerCountOf (int c) { return c & CAPACITY; } private final AtomicInteger ctl = new AtomicInteger (ctlOf(RUNNING, 0 ));
提交优先级
https://camo.githubusercontent.com/fabedab95c31c3d05744f7f68065fc8d22b0b1edb3b0b0addfd869e2f63db36f/68747470733a2f2f63646e2e6e6c61726b2e636f6d2f79757175652f302f323032332f706e672f3331363533332f313637383831313130343436382d62323939636563352d363131312d346638322d626430652d3330303538653034336263632e706e6723617665726167654875653d25323332363237323426636c69656e7449643d7531653063656562392d346330342d342666726f6d3d7061737465266865696768743d3539372669643d753365633234346333266e616d653d696d6167652e706e67266f726967696e4865696768743d353937266f726967696e57696474683d373637266f726967696e616c547970653d62696e61727926726174696f3d3126726f746174696f6e3d302673686f775469746c653d66616c73652673697a653d333435303736267374617475733d646f6e65267374796c653d6e6f6e65267461736b49643d7561343063393939652d666164612d343963662d383866662d3465363862313534343066267469746c653d2677696474683d373637
执行优先级
线程池参数
corePoolSize: 核心线程数,线程池事先创建的线程,当有任务进来的时候第一时间去执行任务,空闲了也不会被回收,会一直重复这些线程
默认情况下,即使是核⼼线程也只能在新任务到达时才创建和启动。但是我们可以使⽤ prestartCoreThread(启动⼀个核⼼线程)或prestartAllCoreThreads(启动全部核⼼线程)⽅法来提前启动 核⼼线程
maximumPoolSize: 最大线程数,当核心线程数处于繁忙并且队列满了的时候,会向操作系统申请额外的线程来消费新进来的任务,允许在这个池子里面最大的线程数量
keepAliveTime: 除了核心线程之外的线程,当执行完任务后,存活空闲下来的时间,超过这个时间就会被回收
unit: 空闲时间的单位,默认是秒
workQueue:工作队列,阻塞队列.分为有界/无界,工作完的线程(不管是core,还是max线程)都会去拉取队列的任务.
threadFactory:线程工厂,设置了一些标识,比如名字
RejectedExecutionHandler:拒绝策略
默认为Abort, 直接抛出异常(AbortPolicy)
什么都不做,也不抛出一场(CallerRunsPolicy)
抛出头部的任务,加入新提交的任务(DiscardOldestPolicy)
谁提交的谁执行(DiscardPolicy)
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 .corePoolSize = corePoolSize; this .maximumPoolSize = maximumPoolSize; this .workQueue = workQueue; this .keepAliveTime = unit.toNanos(keepAliveTime); this .threadFactory = threadFactory; this .handler = handler; }
线程池工作流程
如果当前工作线程数量小于核心线程数量,执行器总是优先 创建一个任务线程,而不是从线程队列中获取一个空闲线程
如果线程池中总的任务数量大于核心线程池数量,新接受的任务将会加入阻塞队列中,一直到阻塞队列已满。在核心线程池数量已经用完,阻塞队列没有满的场景下,线程池不会为新任务创建一个新线程
当完成一个任务执行完时,执行器总是优先从阻塞队列中获取下一个任务,并开始执行,一直到阻塞队列为空,其中所有的缓存任务被取光
在核心线程池数量已经用完、阻塞队列也已经满了的场景下,如果线程池接收到新的任务,将会为新任务创建一个线程(非核心线程),并且立即开始执行新任务
在核心线程都用完、阻塞队列已满的情况下,一直会创建新线程去执行新任务,直到池内的线程总数超出maximumPoolSize。如果线程池的线程总数超过maximumPoolSize,线程池就会拒绝接收任务, 当新任务过来时,会为新任务执行拒绝策略
https://camo.githubusercontent.com/315536b2cbc02ed5c63b3394f45439759d15037c07aa6c51f65358575036ea91/68747470733a2f2f63646e2e6e6c61726b2e636f6d2f79757175652f302f323032332f706e672f3331363533332f313637383834393038313036372d32633562326532332d306335302d346235352d623264312d6363633163643538646336312e706e6723617665726167654875653d25323366626637663526636c69656e7449643d7534643133313834352d303334352d342666726f6d3d7061737465266865696768743d3434372669643d753537616438306262266e616d653d696d6167652e706e67266f726967696e4865696768743d353539266f726967696e57696474683d31323931266f726967696e616c547970653d62696e61727926726174696f3d312e323526726f746174696f6e3d302673686f775469746c653d66616c73652673697a653d3533303039267374617475733d646f6e65267374796c653d6e6f6e65267461736b49643d7562333331366636362d363931662d343062342d623838662d6135336665306566346237267469746c653d2677696474683d313033322e38
https://camo.githubusercontent.com/21cbc5fdb98a9c81a2f1224f8f464f6cffb125d0dbbe8e51c8adde9549341bbf/68747470733a2f2f63646e2e6e6c61726b2e636f6d2f79757175652f302f323032332f706e672f3331363533332f313637383831303835323130312d65353431383238372d626464352d343162342d623566302d3131333865366339323836342e706e6723617665726167654875653d25323366336633656626636c69656e7449643d7531653063656562392d346330342d342666726f6d3d7061737465266865696768743d3334322669643d753732646435313635266e616d653d696d6167652e706e67266f726967696e4865696768743d333432266f726967696e57696474683d373330266f726967696e616c547970653d62696e61727926726174696f3d3126726f746174696f6e3d302673686f775469746c653d66616c73652673697a653d3830333436267374617475733d646f6e65267374796c653d6e6f6e65267461736b49643d7535353963356639372d333033662d343535382d393334662d3962653430373430336563267469746c653d2677696474683d373330
ThreadPoolExecutor执行任务 有 submit 和 execute 二种方式
submit 方法 使用 FutureTask 包裹一层 Runnable,返回 FutureTask 回去,中间调用 execute 方法
execute方法
如果当前正在执行的worker数量比corePoolSize小,直接创建一个新的worker执行任务,调用addWorker方法
如果当前正在执行的worker数量大于等于corePoolSize,将任务放到阻塞队列里,等待空闲线程来执行
若队列的任务数达到上限,且当前运行线程数小于 maximumPoolSize ,任务入队列失败,新创建worker执行任务
若创建线程也失败(队列任务达到上限 且 当前线程数达到了 maximumPoolSize),对于新加入的任务,就会调用reject进行拒绝策略
public Future<?> submit(Runnable task) { if (task == null ) throw new NullPointerException (); RunnableFuture<Void> ftask = newTaskFor(task, null ); execute(ftask); return ftask; } public void execute (Runnable command) { if (command == null ) throw new NullPointerException (); int c = ctl.get(); if (workerCountOf(c) < corePoolSize) { if (addWorker(command, true )) return ; c = ctl.get(); } if (isRunning(c) && workQueue.offer(command)) { int recheck = ctl.get(); if (! isRunning(recheck) && remove(command)) reject(command); else if (workerCountOf(recheck) == 0 ) addWorker(null , false ); else : 仅仅加入阻塞队列中 } else if (!addWorker(command, false )) reject(command); }
addWorker方法 private boolean addWorker (Runnable firstTask, boolean core) { retry: for (;;) { int c = ctl.get(); int rs = runStateOf(c); if (rs >= SHUTDOWN && ! (rs == SHUTDOWN && firstTask == null && ! workQueue.isEmpty())) return false ; for (;;) { int wc = workerCountOf(c); if (wc >= CAPACITY || wc >= (core ? corePoolSize : maximumPoolSize)) return false ; if (compareAndIncrementWorkerCount(c)) break retry; c = ctl.get(); if (runStateOf(c) != rs) continue retry; } } boolean workerStarted = false ; boolean workerAdded = false ; Worker w = null ; try { w = new Worker (firstTask); final Thread t = w.thread; if (t != null ) { final ReentrantLock mainLock = this .mainLock; mainLock.lock(); try { int rs = runStateOf(ctl.get()); if (rs < SHUTDOWN || (rs == SHUTDOWN && firstTask == null )) { if (t.isAlive()) throw new IllegalThreadStateException (); workers.add(w); int s = workers.size(); if (s > largestPoolSize) largestPoolSize = s; workerAdded = true ; } } finally { mainLock.unlock(); } if (workerAdded) { t.start(); workerStarted = true ; } } } finally { if (! workerStarted) addWorkerFailed(w); } return workerStarted; }
runWorker方法 final void runWorker (Worker w) { Thread wt = Thread.currentThread(); Runnable task = w.firstTask; 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 { 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 ; w.completedTasks++; w.unlock(); } } completedAbruptly = false ; } finally { processWorkerExit(w, completedAbruptly); } }
getTask方法 如果发生了下面事中的一个,那么worker需要被回收:
worker个数比线程池最大的还要大
线程池处于STOP状态
线程池处于SHUTDOWN状态并且阻塞队列为空
使用超时时间从阻塞队列里拿数据,并且超时之后没有拿到数据(allowCoreThreadTimeOut || wc > corePoolSize)
如果 getTask 返回的是null,那说明阻塞队列已经没有任务并且当前调用getTask的Worker需要被回收,那么会调用processWorkerExit方法进行回收
private Runnable getTask () { boolean timedOut = false ; for (;;) { int c = ctl.get(); int rs = runStateOf(c); if (rs >= SHUTDOWN && (rs >= STOP || workQueue.isEmpty())) { decrementWorkerCount(); return null ; } int wc = workerCountOf(c); boolean timed = allowCoreThreadTimeOut || wc > corePoolSize; if ((wc > maximumPoolSize || (timed && timedOut)) && (wc > 1 || workQueue.isEmpty())) { if (compareAndDecrementWorkerCount(c)) return null ; continue ; } try { Runnable r = timed ? workQueue.poll(keepAliveTime, TimeUnit.NANOSECONDS) : workQueue.take(); if (r != null ) return r; timedOut = true ; } catch (InterruptedException retry) { timedOut = false ; } } }
processWorkerExit方法 private void processWorkerExit (Worker w, boolean completedAbruptly) { if (completedAbruptly) decrementWorkerCount(); final ReentrantLock mainLock = this .mainLock; mainLock.lock(); try { completedTaskCount += w.completedTasks; 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 ; } addWorker(null , false ); } }
tryTerminate方法 final void tryTerminate () { for (;;) { int c = ctl.get(); if (isRunning(c) || runStateAtLeast(c, TIDYING) || (runStateOf(c) == SHUTDOWN && ! workQueue.isEmpty())) return ; if (workerCountOf(c) != 0 ) { interruptIdleWorkers(ONLY_ONE); return ; } final ReentrantLock mainLock = this .mainLock; mainLock.lock(); try { if (ctl.compareAndSet(c, ctlOf(TIDYING, 0 ))) { try { terminated(); } finally { ctl.set(ctlOf(TERMINATED, 0 )); termination.signalAll(); } return ; } } finally { mainLock.unlock(); } } }
ThreadPoolExecutor关闭 线程池关闭主要有 shutdown 和 shutdownNow 方法 shutdown方法会更新状态到SHUTDOWN,不会影响阻塞队列里任务的执行,但是不会执行新进来的任务。同时也会回收闲置的worker shutdownNow方法会更新状态到STOP,会影响阻塞队列的任务执行,也不会执行新进来的任务
shutdown方法 shutdown 方法,关闭线程池,关闭之后阻塞队列里的任务不受影响,会继续被worker处理,但是新的任务不会被接受
public void shutdown () { final ReentrantLock mainLock = this .mainLock; mainLock.lock(); try { checkShutdownAccess(); advanceRunState(SHUTDOWN); interruptIdleWorkers(); onShutdown(); } finally { mainLock.unlock(); } tryTerminate(); }
interruptIdleWorkers()方法 private void interruptIdleWorkers () { interruptIdleWorkers(false ); }
interruptIdleWorkers(boolean onlyOne)方法 private void interruptIdleWorkers (boolean onlyOne) { final ReentrantLock mainLock = this .mainLock; mainLock.lock(); try { for (Worker w : workers) { Thread t = w.thread; if (!t.isInterrupted() && w.tryLock()) { try { t.interrupt(); } catch (SecurityException ignore) { } finally { w.unlock(); } } if (onlyOne) break ; } } finally { mainLock.unlock(); } }
shutdownNow方法 public List<Runnable> shutdownNow () { List<Runnable> tasks; final ReentrantLock mainLock = this .mainLock; mainLock.lock(); try { checkShutdownAccess(); advanceRunState(STOP); interruptWorkers(); tasks = drainQueue(); } finally { mainLock.unlock(); } tryTerminate(); return tasks; }
interruptWorkers() 方法 private void interruptWorkers () { final ReentrantLock mainLock = this .mainLock; mainLock.lock(); try { for (Worker w : workers) w.interruptIfStarted(); } finally { mainLock.unlock(); } }
interruptIfStarted()方法 void interruptIfStarted () { Thread t; if (getState() >= 0 && (t = thread) != null && !t.isInterrupted()) { try { t.interrupt(); } catch (SecurityException ignore) { } } }
常用线程池 newSingleThreadPool core=1,单线程,max=1,最大只有一条线程,keepAliveTime=0是不回收的,无界队列,一直复用同一个线程,好处是不用一直创建线程,一条线程一直复用 如果异常了,会new一条新线程,能保证顺序执行
public static ExecutorService newSingleThreadExecutor () { return new FinalizableDelegatedExecutorService (new ThreadPoolExecutor (1 , 1 , 0L , TimeUnit.MILLISECONDS, new LinkedBlockingQueue <Runnable>())); }
newCachedThreadPool core=0,max是Integer的最大值,阻塞队列是一个 SynchronousQueue(容量为0的队列,不能放任何东西,并且是一个阻塞队列,作用阻塞任务,让缓存任务队列无效,所有的任务都是走非核心线程,效率非常高),没有池化思想了,一个任务创建一个线程,回收时间为 60 s,当没任务执行/提交的时候,线程池内无线程
时效性高,新任务进来就会有线程去执行任务
不用管理线程的回收,因为线程池管理了线程的回收(keepAliveTime的时间为60s)
资源损耗严重
public static ExecutorService newCachedThreadPool () { return new ThreadPoolExecutor (0 , Integer.MAX_VALUE, 60L , TimeUnit.SECONDS, new SynchronousQueue <Runnable>()); }
newFixedThreadPool core 和 max 的值都是相等,不会额外的创建新线程,keepAliveTime=0,也不会回收空闲线程,无界队列,任务都由核心线程数去执行,保证线程的复用 特点: 适用于,已知任务数量,但是比较耗时的任务 注意: newFixedThreadPool的阻塞队列大小是没有大小限制的,如果队列堆积数据太多会造成资源消耗
public static ExecutorService newFixedThreadPool (int nThreads) { return new ThreadPoolExecutor (nThreads, nThreads, 0L , TimeUnit.MILLISECONDS, new LinkedBlockingQueue <Runnable>()); }
newScheduledThreadPool 线程池支持定时任务及周期性执行任务,创建一个 corePoolSize 传入,最大线程数为整形的最大数的线程池
public static ScheduledExecutorService newScheduledThreadPool (int corePoolSize) { return new ScheduledThreadPoolExecutor (corePoolSize); } public ScheduledThreadPoolExecutor (int corePoolSize) { super (corePoolSize, Integer.MAX_VALUE, 0 , NANOSECONDS, new DelayedWorkQueue ()); }
newWorkStealingPool 建持有足够线程的线程池来达到快速运算的目的,在内部通过使用多个 队列 来减少各个线程调度产生的竞争。这里所说的有足够的线程 JDK 根据 当前线程的运行需求向操作系统申请足够的线程,以保障线程的快速执行,并很大程度 地使用系统资源,提高并发计算的效率,省去用户根据 CPU 资源估算并行度的过程 然,如果开发者想自己定义线程的并发数,则也可以将其作为参数传入
public static ExecutorService newWorkStealingPool () { return new ForkJoinPool (Runtime.getRuntime().availableProcessors(), ForkJoinPool.defaultForkJoinWorkerThreadFactory, null , true ); }
动态线程池 动态化线程池的核心设计包括如下:
简化线程池的配置:线程池构造参数有8个,但是最核心的是3个,corePoolSize,maximumPoolSize,workQueue,它们最大程度的决定了线程池的任务分配和线程的分配策略。考虑到在实际应用中我们获取并发行的场景
并行执行子任务,提高响应速度。这种情况下应该使用同步队列,没有什么任务应该被缓下来,而是应该立即执行
并行执行大批次任务,提高吞吐量。这种情况下,使用使用有界队列,使用队列取缓冲大批量的任务,队列容量必须声明,防止任务无限制堆积。所以线程池只需要提供这三个关键参数的配置,并且提供两种队列的选择,就可以满足绝大多数的业务需求,Less is More
参数可动态的修改:为了解决参数不好配,修改参数成本高等问题,在Java线程池留有高扩展性的基础上,封装线程池,允许线程池监听同步外部的消息,根据消息进行修改配置。将线程池的配置放置在平台侧,允许开发同学简单的查看、修改线程池配置。
增加线程池监控,对某事物缺乏状态的观测,就对其改进无从下手。在线程池执行任务的生命周期添加监控能力,帮助开发同学了解线程池状态
🤗 总结归纳 📎 参考文章
💡 有关文章的问题,欢迎您在底部评论区留言,一起交流~