Day25 线程池
线程池
-
简介
线程池(Thread Pool)是一种基于池化思想管理线程的工具,线程池维护多个线程,等待监督管理者分配可并发执行的任务。避免了处理任务时创建销毁线程开销的代价,另外也避免了线程数量膨胀导致的过分调度问题,保证了对内核的充分利用。
- 降低资源消耗:通过池化技术重复利用已创建的线程,降低线程创建和销毁造成的损耗。
- 提高响应速度:任务到达时,无需等待线程创建即可立即执行。
- 提高线程的可管理性:线程是稀缺资源,如果无限制创建,不仅会消耗系统资源,还会因为线程的不合理分布导致资源调度失衡,降低系统的稳定性。使用线程池可以进行统一的分配、调优和监控。
- 提供更多更强大的功能:线程池具备可拓展性,允许开发人员向其中增加更多的功能。比如延时定时线程池ScheduledThreadPoolExecutor,就允许任务延期执行或定期执行。
-
线程池使用
-
ThreadPoolExecutor 推荐
-
核心参数
-
corePoolSize:核心线程数,线程池初始化时默认是没有线程的,当任务来临时才开始创建线程去执行任务
-
maximumPoolSize:最大线程数,在核心线程数已满,且队列已满时,如果池子里的工作线程数小于maximumPoolSize,则会创建非核心线程执行任务。一般设置为 (最大任务数-任务队列长度)* 单个任务执行时间
-
keepAliveTime:非核心线程数的空闲时间超过keepAliveTime就会被自动终止回收掉,但在corePoolSize=maximumPoolSize时,该值无效,因为不存在非核心线程
-
unit:keepAliveTime的时间单位
-
workQueue:用于保存线程任务的队列,主要分为无界、有界、同步移交等队列,当池子里的工作线程数大于corePoolSize,就会将新进来的线程任务放入队列中 【见并发容器——BlockQueue】一般设置为核心线程数/单个任务执行耗时 * 2
▼java复制代码// private Runnable getTask() => for(;;) try { // 使用阻塞队列,万一队列为空,那这里就会阻塞住直到有任务 // 如果使用非阻塞队列,这里的死循环会反复空转,浪费CPU资源 Runnable r = timed ? workQueue.poll(keepAliveTime, TimeUnit.NANOSECONDS) : workQueue.take(); if (r != null) return r; timedOut = true; } catch (InterruptedException retry) { timedOut = false; } -
threadFactory:
- 创建线程的工厂接口,默认使用Executors.defaultThreadFactory()
- 另外可以实现ThreadFactory接口,自定义线程工厂
-
handler:线程池无法继续接收任务时(workQueue已满和maximumPoolSize已满)的拒绝策略
- AbortPolicy:默认拒绝策略,线程池和队列都满时不再接受新的提交,抛出RejectedExecutionException异常
- CallerRunsPolicy:让提交任务的主线程来执行任务
- DiscardOldestPolicy:丢弃在队列中存在时间最久的任务,重复执行
- DiscardPolicy:丢弃任务,不进行任何通知
- 另外可以实现RejectedExecutionHandler接口,自定义拒绝策略
-
-
使用示例
▼java复制代码import java.util.concurrent.ArrayBlockingQueue; import java.util.concurrent.ThreadPoolExecutor; import java.util.concurrent.TimeUnit; import java.util.concurrent.RejectedExecutionHandler; import java.util.concurrent.ThreadPoolExecutor.AbortPolicy; public class ThreadPoolExample2 { public static void main(String[] args) { // 线程池的核心线程数 int corePoolSize = 5; // 线程池的最大线程数 int maximumPoolSize = 10; // 线程池的任务队列 ArrayBlockingQueue<Runnable> workQueue = new ArrayBlockingQueue<>(100); // 线程池保持空闲的时间 long keepAliveTime = 60L; // 时间单位 TimeUnit unit = TimeUnit.SECONDS; // 线程池的拒绝策略 RejectedExecutionHandler handler = new ThreadPoolExecutor.AbortPolicy(); // 创建线程池 ThreadPoolExecutor executor = new ThreadPoolExecutor( corePoolSize, maximumPoolSize, keepAliveTime, unit, workQueue, handler ); // 提交任务到线程池 for (int i = 0; i < 15; i++) { final int taskId = i; executor.execute(() -> { System.out.println("Task " + taskId + " is running by " + Thread.currentThread().getName()); }); } // 关闭线程池 executor.shutdown(); try { // 等待所有任务完成,超时时间为60秒 if (!executor.awaitTermination(60, TimeUnit.SECONDS)) { // 如果超时后任务仍未完成,则强制关闭线程池 executor.shutdownNow(); } } catch (InterruptedException e) { // 如果等待过程中被中断,也强制关闭线程池 executor.shutdownNow(); } System.out.println("All tasks are done or interrupted."); } }
-
-
Executors 不推荐
提供一些静态类创建线程池,但风险高,一些关键参数不可控
-
常用线程池类型
- newFixedThreadPool 定长线程池,核心线程池固定,使用无界队列,任务数大于线程数时入队,永不拒绝;会**有OOM风险。**仅适合任务量稳定可预测的小型应用
- newCachedThreadPool 弹性伸缩线程池,使用无上限的非核心线程池,来一个任务就建一个线程;高并发时线程数过多容易导致系统崩溃。仅适合执行大量耗时极短且并发可控的任务。
- newSingleThreadExecutor 定时线程池,支持周期性任务执行,核心线程池固定,非核心线程池无上限,使用无界延迟队列;**线程数爆炸+队列堆积→OOM,**可手动设置maximumPoolsize
- newSingleThreadScheduledExecutor 单线程线程池,只有一个工作线程,任务按顺序执行,使用无界队列,有OOM风险;适合需要异步串行执行任务的场景。
- newWorkStealingPool(JDK8+)基于ForkJoinPool的抢占式线程池,空闲线程会从其他繁忙线程的队列拿任务;适合分治任务
-
使用示例
▼java复制代码import java.util.concurrent.ExecutorService; import java.util.concurrent.Executors; public class ThreadPoolExample { public static void main(String[] args) { // 创建一个固定大小的线程池 ExecutorService executor = Executors.newFixedThreadPool(5); // 提交任务到线程池 for (int i = 0; i < 10; i++) { final int taskId = i; executor.submit(() -> { System.out.println("Task " + taskId + " is running by " + Thread.currentThread().getName()); }); } // 关闭线程池 executor.shutdown(); while (!executor.isTerminated()) { // 等待所有任务完成 } System.out.println("All tasks are done."); } }
-
-
任务提交
-
无返回值的任务使用 public void execute(Runnable command) 方法提交;
-
有返回值的任务使用:
- Future submit(Runnable task) : 提交Runnable任务
- Future submit(Runnable task, T result): 提交Runnable任务并指定执行结果
- Future submit(Callable task) : 提交Callable任务
-
批量任务
▼java复制代码# 执行批量任务,返回它们的执行结果 public <T> List<Future<T>> invokeAll(Collection<? extends Callable<T>> tasks) throws InterruptedException #执行批量任务,返回指定时间内完成的执行结果,取消未完成的任务 public <T> List<Future<T>> invokeAll(Collection<? extends Callable<T>> tasks,long timeout, TimeUnit unit) # 执行批量任务,返回最先完成的执行结果 public <T> T invokeAny(Collection<? extends Callable<T>> tasks) throws InterruptedException, ExecutionException #执行批量任务,返回指定时间内最先完成的执行结果,取消未完成的任务 public <T> T invokeAny(Collection<? extends Callable<T>> tasks,long timeout, TimeUnit unit) -
定时/延时任务
▼java复制代码#具备执行定时、延时、周期性任务的线程池 public class ScheduledThreadPoolExecutor extends ThreadPoolExecutor implements ScheduledExecutorService { #延时执行Runnable任务,只执行一次 public ScheduledFuture<?> schedule(Runnable command,long delay,TimeUnit unit) #延时执行Callable任务,只执行一次 public <V> ScheduledFuture<V> schedule(Callable<V> callable,long delay,TimeUnit unit) #廷时一段时间后,周期性执行Runnable任务,周期为固定时间 public ScheduledFuture<?> scheduleAtFixedRate(Runnable command,long initialDelay,long period,TimeUnit unit) #廷时一段时间后,周期性执行Runnable任务,周期为间隔时间 public ScheduledFuture<?> scheduleWithFixedDelay(Runnable command,long initialDelay,long delay,TimeUnit unit)
-
-
线程池关闭
- shutdownNow():立即关闭线程池,正在执行中的任务和队列中的任务都会被中断,关闭后状态为STOP,同时返回被中断的队列中的任务列表。
- shutdown():延时关闭线程池,正在执行中的任务和队列中的任务都能执行完成,关闭后状态为SHUTDOWN,后续进来的新任务会被执行拒绝策略。
- isTerminated():当正在执行的任务和队列中的任务全部都执行完时返回true。
-
-
原理分析
-
任务执行流程 submit方法底层实际调用的是execute,只是用Future包装了一层
▼java复制代码public Future<?> submit(Runnable task) { if (task == null) throw new NullPointerException(); RunnableFuture<Void> ftask = newTaskFor(task, null); execute(ftask); return ftask; }execute()执行分为三步

在实际运行时,不管线程池中的线程是否空闲,只要数量小于核心线程数就会创建新线程。
线程发生异常时,当前线程会被移出线程池,源码中任务出现异常时,在finally块里会执行processWorkerExit()方法,处理当前线程并新增一个线程,以维持固定的核心线程数
-
线程池的流转状态
- RUNNING:会接收新任务并且会处理队列中的任务
- SHUTDOWN:不会接收新任务并且会处理队列中的任务
- STOP:不会接收新任务并且不会处理队列中的任务,并且会中断在处理的任务(注意:一个任务能不能被中断得看任务本身)
- TIDYING:所有任务都终止了,线程池中也没有线程了,这样线程池的状态就会转为TIDYING,一旦达到此状态,就会调用线程池的terminated()
- TERMINATED:terminated()执行完之后就会转变为TERMINATED
RUNNING → SHUTDOWN:调用shutdown()时触发(GC时会调用)
RUNNING/SHUTDOWN→STOP:调用shutdownNow()时触发
SHUTDOWN→TIDYING:队列为空且线程池中没有线程时自动转换
STOP→TIDYING:线程池中没有线程时转换
TIDYING→TERMINATED:调用terminated()

-
-
线程池参数动态配置
背景: 在日常项目中使用线程池来处理一些并发场景,提高任务处理的效率。但是实际使用时,无法准确地设置线程池参数,只能在运行过程中,不断去调整参数,然后重启服务。

-
基于Nacos实现
借助Nacos的Listener,在Bean初始化的时候,开启Nacos配置变更监听
▼java复制代码@Configuration @Data public class MyDynamicThreadPool implements InitializingBean { @Value("${threadPool.corePoolSize}") private int corePoolSize; @Value("${threadPool.maxPoolSize}") private int maxPoolSize; @Value("${threadPool.queueCapacity}") private int queueCapacity; @Value("${threadPool.keepAliveSeconds}") private int keepAliveSeconds; private static ThreadPoolTaskExecutor threadPoolTaskExecutor; @Autowired private NacosConfigManager nacosConfigManager; @Autowired private NacosConfigProperties nacosConfigProperties; @Override public void afterPropertiesSet() throws Exception { threadPoolTaskExecutor = new ThreadPoolTaskExecutor(); threadPoolTaskExecutor.setCorePoolSize(corePoolSize); threadPoolTaskExecutor.setMaxPoolSize(maxPoolSize); threadPoolTaskExecutor.setQueueCapacity(queueCapacity); threadPoolTaskExecutor.setKeepAliveSeconds(keepAliveSeconds); threadPoolTaskExecutor.setThreadNamePrefix( "Fox--"); threadPoolTaskExecutor.setRejectedExecutionHandler( new RejectedExecutionHandler() { @Override public void rejectedExecution(Runnable r, ThreadPoolExecutor executor) { System.out.println("队列已满,丢弃任务"); } }); threadPoolTaskExecutor.initialize(); nacosConfigManager.getConfigService().addListener("threadPool.yml", nacosConfigProperties.getGroup(), new Listener() { @Override public Executor getExecutor() { return null; } @Override public void receiveConfigInfo(String configInfo) { System.out.println("动态修改前-->"); print(); Yaml yaml = new Yaml(); InputStream inputStream = new ByteArrayInputStream(configInfo.getBytes()); Map<String, Object> dataMap = yaml.load(inputStream); // 将Map转换为JSONObject JSONObject pool = new JSONObject(dataMap).getJSONObject("threadPool"); threadPoolTaskExecutor.setCorePoolSize(pool.getInteger("corePoolSize")); threadPoolTaskExecutor.setMaxPoolSize(pool.getInteger("maxPoolSize")); threadPoolTaskExecutor.setQueueCapacity(pool.getInteger("keepAliveSeconds")); threadPoolTaskExecutor.setQueueCapacity(pool.getInteger("queueCapacity")); System.out.println("动态修改后-->"); print(); } }); } //执行任务 public void execute(Runnable runnable){ threadPoolTaskExecutor.execute(runnable); } public void print(){ System.out.println("核心线程数:" + threadPoolTaskExecutor.getThreadPoolExecutor().getCorePoolSize() + " " +"最大线程数:" + threadPoolTaskExecutor.getThreadPoolExecutor().getMaximumPoolSize() +" " + "阻塞队列数:" + threadPoolTaskExecutor.getThreadPoolExecutor().getQueue().size() + "/" + queueCapacity +" " + "活跃线程数:" + threadPoolTaskExecutor.getThreadPoolExecutor().getActiveCount()); } } -
使用美团开源的DynamicTp-基于配置中心的轻量级动态可监控线程池
支持动态调参、通知告警、运行监控、三方包线程池管理等功能https://dynamictp.cn
▼yaml复制代码spring: dynamic: tp: enabled: true # 是否启用 dynamictp,默认true executors: # 动态线程池配置,都有默认值,采用默认值的可以不配置该项,减少配置量 - threadPoolName: dtpExecutor1 # 线程池名称,必填 threadPoolAliasName: 测试线程池 # 线程池别名,可选 executorType: common # 线程池类型 common、eager、ordered、scheduled、priority,默认 common corePoolSize: 5 # 核心线程数,默认1 maximumPoolSize: 8 # 最大线程数,默认cpu核数 queueCapacity: 2000 # 队列容量,默认1024 queueType: VariableLinkedBlockingQueue # 任务队列,查看源码QueueTypeEnum枚举类,默认VariableLinkedBlockingQueue rejectedHandlerType: CallerRunsPolicy # 拒绝策略,查看RejectedTypeEnum枚举类,默认AbortPolicy keepAliveTime: 10 # 空闲线程等待超时时间,默认60 threadNamePrefix: Fox # 线程名前缀,默认dtp allowCoreThreadTimeOut: false # 是否允许核心线程池超时,默认false waitForTasksToCompleteOnShutdown: true # 参考spring线程池设计,优雅关闭线程池,默认true awaitTerminationSeconds: 5 # 优雅关闭线程池时,阻塞等待线程池中任务执行时间,默认3,单位(s) preStartAllCoreThreads: false # 是否预热所有核心线程,默认false
-




