十年匠心定制 · 商业建站与技术教学双线并行 咨询热线:400-886-1026 service@lmnt.cn
ARTICLE DETAIL

资讯详情

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

从源码看透 Java 线程池:Executor 组件结构与 ThreadPoolExecutor 执行原理

从源码看透 Java 线程池:Executor 组件结构与 ThreadPoolExecutor 执行原理 从源码看透 Java 线程池Executor 组件结构与 ThreadPoolExecutor 执行原理【免费下载链接】source-code-hunter 从源码层面剖析挖掘互联网行业主流技术的底层实现原理为广大开发者 “提升技术深度” 提供便利。目前开放 Spring 全家桶Mybatis、Netty、Dubbo 框架及 Redis、Tomcat 中间件等项目地址: https://gitcode.com/GitHub_Trending/so/source-code-hunter导读线程池是 Java 并发编程中复用线程、控制资源开销的核心组件也是 Spring、Netty、Dubbo、Tomcat 等主流框架处理异步任务的共同基石。本文以 Executor线程池组件.md 为骨架从 Executor、ExecutorService 接口到 AbstractExecutorService 抽象类、ThreadPoolExecutor 实现类再到 Executors 工具类逐层拆解线程池的类结构、核心参数、任务提交与执行流程并结合仓库中 Java并发编程在各主流框架中的应用.md 的实战讲解帮助读者真正掌握线程池的设计思想与调优方法。线程池核心组件图解看源码之前先了解一下该组件最主要的几个接口、抽象类和实现类的结构关系类图如下。该组件中Executor和ExecutorService接口定义了线程池最核心的几个方法提交任务submit()、关闭线程池shutdown()。抽象类AbstractExecutorService主要对公共行为submit()系列方法进行了实现这些submit()方法的实现使用了模板方法模式其中调用的execute()方法是未实现的、来自Executor接口的方法。实现类ThreadPoolExecutor则对线程池进行了具体而复杂的实现。另外还有一个常见的工具类Executors里面为开发者封装了一些可以直接拿来用的线程池。源码赏析Executor 和 ExecutorService 接口线程池的最顶层是Executor接口它只定义了一个方法execute(Runnable)用于在将来的某个时间执行给定的任务该任务可以在新线程、池线程或调用线程中执行。ExecutorService接口继承了Executor在它的基础上扩展出了线程池完整的能力边界优雅关闭、提交带返回值的任务等。public interface Executor { /** * 在将来的某个时间执行给定的 Runnable。该 Runnable 可以在新线程、池线程或调用线程中执行。 */ void execute(Runnable command); } public interface ExecutorService extends Executor { /** * 优雅关闭该关闭会继续执行完以前提交的任务但不再接受新任务。 */ void shutdown(); /** * 提交一个有返回值的任务并返回该任务的未来执行完成后的结果。 * Future的 get()方法 将在成功完成后返回任务的结果。 */ T FutureT submit(CallableT task); T FutureT submit(Runnable task, T result); Future? submit(Runnable task); }从接口设计上可以看出线程池的三大核心职责执行任务execute、提交并获取异步结果submit Future、生命周期管理shutdown。submit()系列方法之所以有多个重载是为了兼容三种任务形态Runnable无返回值、Runnable 指定返回结果、CallableT有返回值。AbstractExecutorService 抽象类AbstractExecutorService实现了ExecutorService中的submit()系列方法它是“模板方法模式”的典型示范将公共的、与具体线程策略无关的逻辑把任务包装成 Future、调用 execute放在抽象类中实现而把最关键的可变点execute()留给子类去定制。/** * 该抽象类最主要的内容就是实现了 ExecutorService 中的 submit()系列方法 */ public abstract class AbstractExecutorService implements ExecutorService { /** * 提交任务 进行执行返回获取未来结果的 Future对象。 * 这里使用了 模板方法模式execute()方法来自 Executor接口该抽象类中并未进行实现 * 而是交由子类具体实现。 */ public Future? submit(Runnable task) { if (task null) throw new NullPointerException(); RunnableFutureVoid ftask newTaskFor(task, null); execute(ftask); return ftask; } public T FutureT submit(Runnable task, T result) { if (task null) throw new NullPointerException(); RunnableFutureT ftask newTaskFor(task, result); execute(ftask); return ftask; } public T FutureT submit(CallableT task) { if (task null) throw new NullPointerException(); RunnableFutureT ftask newTaskFor(task); execute(ftask); return ftask; } }三个重载的实现逻辑高度一致先用newTaskFor()把任务包装成RunnableFuture内部由FutureTask承载任务执行完成后结果可以被 Future 获取再调用模板方法execute(ftask)真正提交执行最后把 Future 返回给调用方调用方通过future.get()阻塞等待并取得任务结果。ThreadPoolExecutor核心参数与构造方法ThreadPoolExecutor继承自AbstractExecutorService是线程池最核心、最复杂的实现类。它的核心字段如下。public class ThreadPoolExecutor extends AbstractExecutorService { /** 阻塞队列 */ private final BlockingQueueRunnable workQueue; /** 用于创建线程的 线程工厂 */ private volatile ThreadFactory threadFactory; /** 核心线程数 */ private volatile int corePoolSize; /** 最大线程数 */ private volatile int maximumPoolSize; ... }各参数的核心语义结合 Java并发编程在各主流框架中的应用.md 中更完整的注释参数类型语义corePoolSizeint核心线程数。提交任务时若已创建线程数小于corePoolSize即使存在空闲线程也会新建线程执行任务直到线程数大于等于corePoolSizemaximumPoolSizeint最大线程数。当队列满了、且已创建线程数小于maximumPoolSize时线程池会创建新线程来执行任务对于无界队列可忽略该参数keepAliveTimelong线程存活保持时间。当线程数超出核心线程数且线程空闲时间超过keepAliveTime该线程会被销毁直到线程数小于等于核心线程数workQueueBlockingQueueRunnable任务队列。用于传输和保存等待执行任务的阻塞队列threadFactoryThreadFactory线程工厂。用于创建新线程默认工厂创建的线程名具有统一风格pool-m-thread-nm 为线程池编号n 为线程池内线程编号handlerRejectedExecutionHandler饱和策略。当线程池和队列都满了再提交的任务按此策略处理ThreadPoolExecutor提供了多个重载构造方法但最终都汇聚到参数最全的那个构造方法上完成初始化并在其中完成参数合法性校验。public ThreadPoolExecutor(int corePoolSize, int maximumPoolSize, long keepAliveTime, TimeUnit unit, BlockingQueueRunnable workQueue) { this(corePoolSize, maximumPoolSize, keepAliveTime, unit, workQueue, Executors.defaultThreadFactory(), defaultHandler); } public ThreadPoolExecutor(int corePoolSize, int maximumPoolSize, long keepAliveTime, TimeUnit unit, BlockingQueueRunnable workQueue, ThreadFactory threadFactory) { this(corePoolSize, maximumPoolSize, keepAliveTime, unit, workQueue, threadFactory, defaultHandler); } public ThreadPoolExecutor(int corePoolSize, int maximumPoolSize, long keepAliveTime, TimeUnit unit, BlockingQueueRunnable workQueue, RejectedExecutionHandler handler) { this(corePoolSize, maximumPoolSize, keepAliveTime, unit, workQueue, Executors.defaultThreadFactory(), handler); } public ThreadPoolExecutor(int corePoolSize, int maximumPoolSize, long keepAliveTime, TimeUnit unit, BlockingQueueRunnable 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; }值得注意的校验细节可直接作为自建线程池时的正确性参考corePoolSize不能为负数maximumPoolSize必须大于 0且maximumPoolSize不能小于corePoolSizekeepAliveTime不能为负数workQueue、threadFactory、handler均不能为 nullkeepAliveTime会统一通过unit.toNanos(keepAliveTime)转换为纳秒存储不传threadFactory与handler时使用默认实现Executors.defaultThreadFactory()与defaultHandler即AbortPolicy。execute()任务提交的三步核心流程execute(Runnable command)是线程池执行任务的入口其判断逻辑可以凝练为“核心线程数 → 工作队列 → 最大线程数”三层递进判断/** 执行 Runnable任务 */ public void execute(Runnable command) { if (command null) throw new NullPointerException(); /* * 分三步进行 * * 1、如果运行的线程少于 corePoolSize尝试开启一个新的线程否则尝试进入工作队列 * * 2. 如果工作队列没满则进入工作队列否则 判断是否超出最大线程数 * * 3. 如果未超出最大线程数则尝试开启一个新的线程否则 按饱和策略处理无法执行的任务 */ 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 if (!addWorker(command, false)) reject(command); }三步逻辑逐条展开第一步核心线程数判断。若当前工作线程数workerCountOf(c)小于corePoolSize则调用addWorker(command, true)直接创建新线程执行任务第二个参数true表示按核心线程数限制校验第二步入队判断。若线程数已达核心数且线程池处于运行态isRunning(c)则尝试workQueue.offer(command)将任务放入阻塞队列等待执行。这里有一段经典的**二次检查recheck**逻辑入队成功后再次读取ctl若线程池已被关闭则把任务从队列移除并走拒绝策略若工作线程数已降为 0例如核心线程被回收则addWorker(null, false)补建一个空任务线程去消费队列第三步最大线程数判断。若入队失败队列已满且addWorker(command, false)第二个参数false表示按最大线程数限制校验也无法创建新线程说明线程数与队列容量都已耗尽调用reject(command)按饱和策略处理。整个流程可以用下图直观表示。这里需要特别强调addWorker(command, true/false)中第二个布尔参数的作用为true时受corePoolSize约束为false时受maximumPoolSize约束这正是“核心线程先扩容、队列满了才扩到最大线程”这一调度语义的实现关键。shutdown() 与生命周期管理shutdown()是优雅关闭不再接受新任务但会继续执行完已经提交的任务如果线程池已经关闭再次调用不会有其他效果。public void shutdown() { final ReentrantLock mainLock this.mainLock; mainLock.lock(); try { checkShutdownAccess(); advanceRunState(SHUTDOWN); interruptIdleWorkers(); onShutdown(); // hook for ScheduledThreadPoolExecutor } finally { mainLock.unlock(); } tryTerminate(); }实现要点通过mainLockReentrantLock保证关闭操作的线程安全checkShutdownAccess()校验安全管理权限advanceRunState(SHUTDOWN)把运行状态原子地推进为SHUTDOWNinterruptIdleWorkers()中断空闲的工作线程让它们退出等待onShutdown()是留给ScheduledThreadPoolExecutor的钩子方法最后tryTerminate()尝试终止线程池。与shutdown()对应的是shutdownNow()它把状态推进到STOP中断所有工作线程通过drainQueue()把队列中未执行的任务返回给调用方这些任务会被抛弃不再执行实现“立即关闭”。判断线程池是否已关闭可用isShutdown()其实现为!isRunning(ctl.get())。从这里可以看出线程池的运行状态与工作线程数被打包在同一个原子变量ctl中通过位运算拆分这也是线程池高效线程安全的底层设计之一。工具类 Executors看类名也知道它最主要的作用就是提供static的工具方法为开发者提供各种封装好的、具有各自特性的线程池。public class Executors { /** * 创建一个固定线程数量的线程池 */ public static ExecutorService newFixedThreadPool(int nThreads) { return new ThreadPoolExecutor(nThreads, nThreads, 0L, TimeUnit.MILLISECONDS, new LinkedBlockingQueueRunnable()); } /** * 创建一个单线程的线程池 */ public static ExecutorService newSingleThreadExecutor() { return new FinalizableDelegatedExecutorService (new ThreadPoolExecutor(1, 1, 0L, TimeUnit.MILLISECONDS, new LinkedBlockingQueueRunnable())); } /** * 创建一个缓存的可动态伸缩的线程池。 * 可以看出来核心线程数为0最大线程数为Integer.MAX_VALUE如果任务数在某一瞬间暴涨 * 这个线程池很可能会把 服务器撑爆。 * 另外需要注意的是它们底层都是使用了 ThreadPoolExecutor只不过帮我们配好了参数 */ public static ExecutorService newCachedThreadPool() { return new ThreadPoolExecutor(0, Integer.MAX_VALUE, 60L, TimeUnit.SECONDS, new SynchronousQueueRunnable()); } }结合仓库中的扩展说明Executors通过ThreadPoolExecutor封装了 4 种常用的线程池其功能与适用场景如下CachedThreadPool缓存型、几乎可无限扩大的线程池最大线程数为Integer.MAX_VALUE适用于执行大量短生命周期的异步任务FixedThreadPool固定大小线程池线程数可控不会造成线程过多、系统负载严重SingleThreadExecutor单线程线程池可以保证任务按调用顺序执行ScheduledThreadPool适用于执行延时或周期性任务。每种封装背后的参数组合值得玩味newFixedThreadPool核心线程数 最大线程数 nkeepAliveTime为 0配无界队列LinkedBlockingQueue因此线程数恒定、不会回收也不会超扩newSingleThreadExecutor参数与 Fixed 一致n1区别在于外层包了FinalizableDelegatedExecutorService委托类屏蔽了ThreadPoolExecutor暴露出来的可配置方法保证单线程语义不被破坏newCachedThreadPool核心线程数为 0、最大线程数为Integer.MAX_VALUE、空闲 60 秒回收配SynchronousQueue不缓存任务的直接交接队列任务一来就建线程执行、空闲即回收——若任务数在某一瞬间暴涨这个线程池很可能会把服务器撑爆使用时务必评估流量峰值。饱和策略队列与线程都满时的兜底当线程池和队列都已满、无法再接收新任务时会触发RejectedExecutionHandler饱和策略。JDK 内置了四种实现策略行为AbortPolicy默认直接抛出RejectedExecutionExceptionCallerRunsPolicy由提交任务的调用者线程直接执行该任务DiscardPolicy静默丢弃无法执行的任务不抛异常DiscardOldestPolicy丢弃队列中最旧的任务然后尝试重新提交当前任务默认策略是defaultHandler即AbortPolicy。从源码注释可以确认handler字段用volatile修饰、可在运行期通过setRejectedExecutionHandler动态更换这为生产环境的动态调优留出了空间。线程池的配置建议与实际应用如何配置线程池结合 Java并发编程在各主流框架中的应用.md 的实践总结线程池大小应根据任务类型区分设置CPU 密集型任务尽量使用较小的线程池一般为CPU 核心数 1。因为 CPU 密集型任务使 CPU 使用率很高若开过多线程数会造成 CPU 过度切换反而降低吞吐IO 密集型任务可以使用稍大的线程池一般为2 × CPU 核心数。IO 密集型任务 CPU 使用率并不高可以让 CPU 在等待 IO 的时候有其他线程去处理别的任务充分利用 CPU 时间。需要明确的是上述经验值适用于 JDK 1.8 时代以 CPU 核心数为基准的常规估算具体生产环境还应结合任务耗时分布、队列容量、RejectedExecutionHandler的选择做压测验证避免照搬公式。线程池的实际应用线程池并非停留在理论层的组件主流框架的底层都离不开它Tomcat在分发 Web 请求时使用线程池来处理并发连接与请求Netty的EventLoop线程模型、Dubbo的远程通信、RocketMQ的消费与拉取流程等都以ThreadPoolExecutor或其变体作为并发基石本仓库中 Dubbo远程通信模块简析.md、EventLoop组件.md、rocketmq-consumer-start.md 等文档均涉及线程池在这些框架中的落地形态可作为延伸阅读。理解Executor组件就等于掌握了阅读这些框架并发代码的通用钥匙看一个框架的线程模型本质上就是看它如何配置corePoolSize、maximumPoolSize、workQueue与handler这四个参数。小结从接口到实现Executor组件呈现出一条清晰的设计主线Executor 接口定义最朴素的任务执行能力execute(Runnable)ExecutorService 接口扩展出提交带返回值任务submitFuture与生命周期管理shutdown能力AbstractExecutorService 抽象类以模板方法模式沉淀公共的submit()系列实现把execute()留给子类ThreadPoolExecutor通过“核心线程数 → 工作队列 → 最大线程数”三步判断完成任务的调度辅以ctl原子状态管理、二次检查、饱和策略兜底实现高效且健壮的线程复用Executors 工具类基于ThreadPoolExecutor参数组合封装出 4 种开箱即用的线程池。掌握这条主线之后再看各类框架源码中千变万化的线程池用法都能迅速定位到它配置了什么队列、什么线程数、什么拒绝策略从而理解其背后的并发模型设计意图。更多并发组件AQS、Lock、Semaphore 等可继续阅读 JUC并发包UML全量类图.md、详解AbstractQueuedSynchronizer.md 与 Lock锁组件.md 系统学习。【免费下载链接】source-code-hunter 从源码层面剖析挖掘互联网行业主流技术的底层实现原理为广大开发者 “提升技术深度” 提供便利。目前开放 Spring 全家桶Mybatis、Netty、Dubbo 框架及 Redis、Tomcat 中间件等项目地址: https://gitcode.com/GitHub_Trending/so/source-code-hunter创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表