執行緒池作用
降低資源消耗:通過池化技術重復利用已創建的執行緒,降低執行緒創建和銷毀造成的損耗,提高回應速度:任務到達時,無需等待執行緒創建即可立即執行,提高執行緒的可管理性:執行緒是稀缺資源,如果無限制創建,不僅會消耗系統資源,還會因為執行緒的不合理分布導致資源調度失衡,降低系統的穩定性,使用執行緒池可以進行統一的分配、調優和監控,提供更多更強大的功能:執行緒池具備可拓展性,允許開發人員向其中增加更多的功能,比如延時定時執行緒池ScheduledThreadPoolExecutor,就允許任務延期執行或定期執行
執行緒池5種狀態
RUNNING:執行緒池一旦被創建,就處于 RUNNING 狀態,任務數為 0,能夠接收新任務,對已排隊的任務進行處理,SHUTDOWN:不接收新任務,但能處理已排隊的任務,呼叫執行緒池的 shutdown() 方法,執行緒池由 RUNNING 轉變為 SHUTDOWN 狀態,STOP:不接收新任務,不處理已排隊的任務,并且會中斷正在處理的任務,呼叫執行緒池的 shutdownNow() 方法,執行緒池由(RUNNING 或 SHUTDOWN ) 轉變為 STOP 狀態,TIDYING:- SHUTDOWN 狀態下,任務數為0, 其他所有任務已終止,執行緒池會變為 TIDYING 狀態,會執行 terminated() 方法,執行緒池中的 terminated() 方法是空實作,可以重寫該方法進行相應的處理,
- 執行緒池在 SHUTDOWN 狀態,任務佇列為空且執行中任務為空,執行緒池就會由 SHUTDOWN 轉變為 TIDYING 狀態,
- 執行緒池在 STOP 狀態,執行緒池中執行中任務為空時,就會由 STOP 轉變為 TIDYING 狀態,
TERMINATED:執行緒池徹底終止,執行緒池在 TIDYING 狀態執行完 terminated() 方法就會由 TIDYING 轉變為 TERMINATED 狀態,

Excutes
newFixedThreadPool- 創建一個固定的長度的執行緒池,每當提交一個任務就創建一個執行緒,知道達到執行緒池的最大數量,這時執行緒規模將不再變化,當執行緒發生未預期的錯誤而結束時,執行緒池會創建一個新的執行緒繼續運行佇列里的任務,
newSingleThreadExecutor- 這是一個單執行緒的 Executor ,它創建單個作業執行緒來執行任務,如果這個執行緒例外結束,會創建一個新的來代替它;它的特點是能確保依照任務在佇列中的順序來串行執行,
newCachedThreadPool- 創建一個可快取的執行緒池,如果執行緒池的規模超過了處理需求,將自動回收空閑執行緒,而當需求正駕駛,則可以自動添加新執行緒,執行緒池的規模不存在任何限制,
newScheduledThreadPool- 創建一個固定長度的執行緒池,而且以延遲或定時的方式來執行任務,類似于 Timer
newSingleThreadScheduledExecutor- 單執行緒可執行周期性任務的執行緒池
newWorkStealingPool- 任務竊取執行緒池,不保證執行順序,適合任務耗時差異較大,執行緒池中有多個執行緒佇列,有的執行緒佇列中有大量的比較耗時的任務堆積,而有的執行緒佇列卻是空的,就存在有的執行緒處于饑餓狀態,當一個執行緒處于饑餓狀態時,它就會去其它的執行緒佇列中竊取任務,解決饑餓導致的效率問題,默認創建的并行 level 是 CPU 的核數,主執行緒結束,即使執行緒池有任務也會立即停止,
為什么不推薦Executors
Executors工具類創建的執行緒池佇列或執行緒默認為Integer.MAX_VALUE,容易堆積請求 阿里巴巴Java開發手冊:
FixedThreadPool和SingleThreadPool:允許的請求佇列長度為 Integer.MAX_VALUE,可能會堆積大量的請求,從而導致 OOMCachedThreadPool:允許的創建執行緒數量為 Integer.MAX_VALUE,可能會創建大量的執行緒,從而導致 OOM
推薦使用ThreadPoolExecutor類根據實際需要自定義創建
ThreadPoolExecutor
七大引數
ThreadPoolExecutor類主要有以下七個引數:
int corePoolSize: 核心執行緒池大小,執行緒池中常駐執行緒的最大數量int maximumPoolSize: 最大核心執行緒池大小(包括核心執行緒和非核心執行緒)long keepAliveTime: 執行緒空閑后的存活時間(僅適用于非核心執行緒)TimeUnit unit: 超時單位BlockingQueue<Runnable> workQueue: 阻塞佇列ThreadFactory threadFactory: 執行緒工廠:創建執行緒的,一般默認RejectedExecutionHandler handle: 拒絕策略
四大策略
拒絕策略就是當佇列滿時,執行緒如何去處理新來的任務,
AbortPolicy(中止策略)-默認
- 功能:當觸發拒絕策略時,直接拋出
拒絕執行的例外 - 使用場景:ThreadPoolExecutor中
默認的策略就是AbortPolicy,由于ExecutorService介面的系列ThreadPoolExecutor都沒有顯示的設定拒絕策略,所以默認的都是這個,
CallerRunsPolicy(呼叫者運行策略)
- 功能:只要執行緒池沒有關閉,就由提交任務的
當前執行緒處理, - 使用場景:一般在不允許失敗、對性能要求不高、并發量較小的場景下使用,
DiscardPolicy(丟棄策略)
- 功能:直接
丟棄這個任務,不觸發任何動作 - 使用場景: 提交的任務無關緊要,一般用的少,
DiscardOldestPolicy(棄老策略)
- 功能:拋棄下一個將要被執行的任務,相當于排隊的時候
把第一個人打死,然后自己代替 - 使用場景:發布訊息、修改訊息類似場景,當老訊息還未執行,此時新的訊息又來了,這時未執行的訊息的版本比現在提交的訊息版本要低就可以被丟棄了,
作業佇列
ArrayBlockingQueue:使用陣列實作的有界阻塞佇列,特性先進先出LinkedBlockingQueue:使用鏈表實作的阻塞佇列,特性先進先出,可以設定其容量,默認為Interger.MAX_VALUEPriorityBlockingQueue:使用平衡二叉樹堆,實作的具有優先級的無界阻塞佇列DelayQueue:無界阻塞延遲佇列,佇列中每個元素均有過期時間,當從佇列獲取元素時,只有過期元素才會出佇列,佇列頭元素是最塊要過期的元素,SynchronousQueue:一個不存盤元素的阻塞佇列,每個插入操作,必須等到另一個執行緒呼叫移除操作,否則插入操作一直處于阻塞狀態
運行流程
- 判斷執行緒池里的核心執行緒是否都在執行任務?
- 否:呼叫/創建一個新的核心執行緒來執行任務
- 是:作業佇列是否已滿?
- 否:將新提交的任務存盤在作業佇列里
- 是:執行緒池里的執行緒數是否達到最大執行緒值?
- 否:呼叫/創建一個新的非核心執行緒來執行任務
- 是:執行執行緒池飽和策略

一個面試題
一個執行緒池 core 7; max 20 , queue: 50, 100并發進來怎么分配的?
答:先有7個能直接得到執行, 接下來把50個進入佇列排隊等候, 在多開13個繼續執行, 現在 70 個被安排上了, 剩下 30 個默認執行飽和策略,
執行任務
execute提交沒有回傳值,不能判斷是否執行成功,只能提交一個Runnable的物件submit會回傳一個Future物件,通過Future的get()方法來獲取回傳值,submit提交執行緒可以吃掉執行緒中產生的例外,達到執行緒復用,當get()執行結果時例外才會拋出,原因是通過submit提交的執行緒,當發生例外時,會將例外保存,待future.get()時才會拋出,
關閉執行緒池
shutdown():不再繼續接收新的任務,執行完成已有任務后關閉shutdownNow():直接關閉,若果有任務嘗試停止
執行緒池出現例外會發生什么?
- 執行緒出現例外,執行緒會退出,并重新創建新的執行緒執行佇列里任務,不能復用執行緒
- 當業務代碼的例外捕獲了,執行緒就可以復用
- 使用ThreadFactory的UncaughtExceptionHandler保證執行緒的所有例外都能捕獲(包括業務的例外),兜底的.如果提交方式用execute,不能復用執行緒
- setUncaughtExceptionHandler+submit :可以吃掉例外并復用執行緒(是吃掉,不報錯)
- setUncaughtExceptionHandler+submit+future.get() :可以獲取到例外并復用執行緒
最佳實踐
- 提交執行緒的
業務例外用try catch處理,保證執行緒不會例外退出 業務之外的例外我們不可預見的,創建執行緒池設定ThreadFactory的UncaughtExceptionHandler可以對未捕獲的例外做保底處理,通過submit提交任務,可以吃掉例外并復用執行緒;想要捕獲例外這時用future.get()
注:關于例外處理的相關案例,已在原始碼中,這里不做展示
實戰1:結合CompletableFuture使用執行緒池
- CompletableFuture,結合了Future的優點,提供了非常強大的Future的擴展功能,可以幫助我們簡化異步編程的復雜性,提供了函式式編程的能力,可以通過回呼的方式處理計算結果,并且提供了轉換和組合CompletableFuture的方法,
- CompletableFuture可以傳入自定義執行緒池,否則使用自己默認的執行緒池,我們習慣做法是自定義執行緒池,控制整個專案的執行緒數量,不使用自定義的執行緒池,做到可控可調
步驟1:宣告一個執行緒池bean
application.properties
//盡量做到每個業務使用自己配置的執行緒池
service1.thread.coreSize=10
service1.thread.maxSize=100
service1.thread.keepAliveTime=10
復制代碼
執行緒池屬性類
/**
* @Description: 執行緒池屬性
* @Author: jianweil
* @date: 2021/12/9 10:44
*/
@ConfigurationProperties(prefix = "service1.thread")
@Data
public class ThreadPoolConfigProperties {
private Integer coreSize;
private Integer maxSize;
private Integer keepAliveTime;
}
復制代碼
執行緒池配置類
/**
* @Description: 執行緒池配置類:根據不同業務定義不同的執行緒池配置
**/
@EnableConfigurationProperties(ThreadPoolConfigProperties.class)
@Configuration
public class MyService1ThreadConfig {
@Bean
public ThreadPoolExecutor threadPoolExecutor(ThreadPoolConfigProperties pool) {
return new ThreadPoolExecutor(
pool.getCoreSize(),
pool.getMaxSize(),
pool.getKeepAliveTime(),
TimeUnit.SECONDS,
new LinkedBlockingDeque<>(100000),
Executors.defaultThreadFactory(),
new ThreadPoolExecutor.AbortPolicy()
);
}
}
復制代碼
步驟2:使用
注:本文所有原始碼已分享github
/**
* @Description: 測驗CompletableFuture
* @Author: jianweil
* @date: 2021/12/9 10:50
*/
@SpringBootTest
public class CompletableFutureTest {
@Autowired
private ThreadPoolExecutor threadPoolExecutor;
/***
* 無回傳值
* runAsync
*/
@Test
public void main1() {
System.out.println("main.................start.....");
CompletableFuture.runAsync(() -> {
System.out.println("當前執行緒:" + Thread.currentThread().getId());
int i = 10 / 2;
System.out.println("運行結果:" + i);
}, threadPoolExecutor);
System.out.println("main.................end......");
}
}
復制代碼
實戰2:結合@Async使用執行緒池
- 在現實的互聯網專案開發中,針對高并發的請求,一般的做法是高并發介面單獨執行緒池隔離處理,可能為某一高并發的介面單獨一個執行緒池
方式1:默認執行緒池
- 使用@Async注解,在默認情況下用的是SimpleAsyncTaskExecutor執行緒池,該執行緒池不是真正意義上的執行緒池
- 使用此執行緒池無法實作執行緒重用,每次呼叫都會新建一條執行緒,若系統中不斷的創建執行緒,最侄訓導致系統占用記憶體過高,引發
OutOfMemoryError錯誤
步驟1:自定義一個能查看執行緒池引數的類
- 不清楚執行緒池當時的情況,有多少執行緒在執行,多少在佇列中等待呢?
- 創建了一個ThreadPoolTaskExecutor的子類,在每次提交執行緒任務的時候都會將當前執行緒池的運行狀況列印出來
public class VisiableThreadPoolTaskExecutor extends ThreadPoolTaskExecutor {
private static final Logger logger = LoggerFactory.getLogger(VisiableThreadPoolTaskExecutor.class);
private void showThreadPoolInfo(String prefix) {
ThreadPoolExecutor threadPoolExecutor = getThreadPoolExecutor();
if (null == threadPoolExecutor) {
return;
}
logger.info("{}, {},taskCount [{}], completedTaskCount [{}], activeCount [{}], queueSize [{}]",
this.getThreadNamePrefix(),
prefix,
threadPoolExecutor.getTaskCount(),
threadPoolExecutor.getCompletedTaskCount(),
threadPoolExecutor.getActiveCount(),
threadPoolExecutor.getQueue().size());
}
@Override
public void execute(Runnable task) {
showThreadPoolInfo("1. do execute");
super.execute(task);
}
@Override
public void execute(Runnable task, long startTimeout) {
showThreadPoolInfo("2. do execute");
super.execute(task, startTimeout);
}
@Override
public Future<?> submit(Runnable task) {
showThreadPoolInfo("1. do submit");
return super.submit(task);
}
@Override
public <T> Future<T> submit(Callable<T> task) {
showThreadPoolInfo("2. do submit");
return super.submit(task);
}
@Override
public ListenableFuture<?> submitListenable(Runnable task) {
showThreadPoolInfo("1. do submitListenable");
return super.submitListenable(task);
}
@Override
public <T> ListenableFuture<T> submitListenable(Callable<T> task) {
showThreadPoolInfo("2. do submitListenable");
return super.submitListenable(task);
}
}
復制代碼
步驟2:實作AsyncConfigurer類
- 要配置默認的執行緒池,要實作
AsyncConfigurer類的兩個方法 - 不需要列印運行狀況的可以使用ThreadPoolTaskExecutor類構建執行緒池
/**
* @Description: 注解@async配置
* @Author: jianweil
* @date: 2021/12/9 11:52
*/
@Slf4j
@EnableAsync
@Configuration
public class AsyncThreadConfig implements AsyncConfigurer {
/**
* 定義@Async默認的執行緒池
* ThreadPoolTaskExecutor不是完全被IOC容器管理的bean,可以在方法上加上@Bean注解交給容器管理,這樣可以將taskExecutor.initialize()方法呼叫去掉,容器會自動呼叫
*
* @return
*/
@Override
public Executor getAsyncExecutor() {
int processors = Runtime.getRuntime().availableProcessors();
//常用的執行器
//ThreadPoolTaskExecutor taskExecutor = new ThreadPoolTaskExecutor();
//可以查看執行緒池引數的自定義執行器
ThreadPoolTaskExecutor taskExecutor = new VisiableThreadPoolTaskExecutor();
//核心執行緒數
taskExecutor.setCorePoolSize(1);
taskExecutor.setMaxPoolSize(2);
//執行緒佇列最大執行緒數,默認:50
taskExecutor.setQueueCapacity(50);
//執行緒名稱前綴
taskExecutor.setThreadNamePrefix("default-ljw-");
taskExecutor.setRejectedExecutionHandler(new ThreadPoolExecutor.CallerRunsPolicy());
//執行初始化(重要)
taskExecutor.initialize();
return taskExecutor;
}
/**
* 異步方法執行的程序中拋出的例外捕獲
*
* @return
*/
@Override
public AsyncUncaughtExceptionHandler getAsyncUncaughtExceptionHandler() {
return (ex, method, params) ->
log.error("執行緒池執行任務發送未知錯誤,執行方法:{}", method.getName(), ex.getMessage());
}
}
復制代碼
步驟3: 使用
- 直接添加注解@Async即可使用到配置的
默認執行緒池
/**
* 默認執行緒池
*/
@Async
public void defaultThread() throws Exception {
long start = System.currentTimeMillis();
Thread.sleep(random.nextInt(1000));
long end = System.currentTimeMillis();
int i = 1 / 0;
log.info("使用默認執行緒池,耗時:" + (end - start) + "毫秒");
}
復制代碼
方式2:指定執行緒池
- 由于業務需要,根據業務不同需要不同的執行緒池
步驟1:宣告一個執行緒池bean
/**
* @Description: 注解@async配置
* @Author: jianweil
* @date: 2021/12/9 11:52
*/
@Slf4j
@EnableAsync
@Configuration
public class AsyncThreadConfig implements AsyncConfigurer {
@Bean("service2Executor")
public Executor service2Executor() {
//Java虛擬機可用的處理器數
int processors = Runtime.getRuntime().availableProcessors();
//定義執行緒池
ThreadPoolTaskExecutor taskExecutor = new ThreadPoolTaskExecutor();
//可以查看執行緒池引數的自定義執行器
//ThreadPoolTaskExecutor taskExecutor = new VisiableThreadPoolTaskExecutor();
//核心執行緒數
taskExecutor.setCorePoolSize(processors);
taskExecutor.setMaxPoolSize(100);
//執行緒佇列最大執行緒數,默認:100
taskExecutor.setQueueCapacity(100);
//執行緒名稱前綴
taskExecutor.setThreadNamePrefix("my-ljw-");
//執行緒池中執行緒最大空閑時間,默認:60,單位:秒
taskExecutor.setKeepAliveSeconds(60);
//核心執行緒是否允許超時,默認:false
taskExecutor.setAllowCoreThreadTimeOut(false);
//IOC容器關閉時是否阻塞等待剩余的任務執行完成,默認:false(必須設定setAwaitTerminationSeconds)
taskExecutor.setWaitForTasksToCompleteOnShutdown(false);
//阻塞IOC容器關閉的時間,默認:10秒(必須設定setWaitForTasksToCompleteOnShutdown)
taskExecutor.setAwaitTerminationSeconds(10);
/**
* 拒絕策略,默認是AbortPolicy
* AbortPolicy:丟棄任務并拋出RejectedExecutionException例外
* DiscardPolicy:丟棄任務但不拋出例外
* DiscardOldestPolicy:丟棄最舊的處理程式,然后重試,如果執行器關閉,這時丟棄任務
* CallerRunsPolicy:執行器執行任務失敗,則在策略回呼方法中執行任務,如果執行器關閉,這時丟棄任務
*/
taskExecutor.setRejectedExecutionHandler(new ThreadPoolExecutor.CallerRunsPolicy());
return taskExecutor;
}
}
復制代碼
步驟2: 使用
- @Async("service2Executor")注解指定使用的執行緒池名稱
/**
* 指定執行緒池service2Executor
*
* @throws Exception
*/
@Async("service2Executor")
public void service2Executor() throws Exception {
long start = System.currentTimeMillis();
Thread.sleep(random.nextInt(1000));
long end = System.currentTimeMillis();
log.info("使用執行緒池service2Executor,耗時:" + (end - start) + "毫秒");
}
復制代碼
注:異步任務回傳值為void,不能獲取的回傳值的
計算執行緒數量
適合框架類
例如netty,dubbo這種底層通訊框架通常會參考進行設定
- IO 密集型(較多): 通常設定為2n+1,其中n為CPU核數
- CPU 密集型(較少): 通常設定為 n+1
實際情況往往復雜的多,并不會按照這個進行設定
IO密集型型別進階演算法
- 對于IO密集型型別的應用:執行緒數 = CPU核心數/(1-阻塞系數)
- 引入了阻塞系數的概念,一般為0.8~0.9之間,
在我們的業務開發中,基本上都是IO密集型,因為往往都會去操作資料庫,訪問redis,es等存盤型組件,涉及到磁盤IO,網路IO,
一個4C8G的機器上部署了一個MQ消費者,在RocketMQ的實作中,消費端也是用一個執行緒池來消費執行緒的,那這個執行緒數要怎么設定呢?
- 如果按照 2n + 1 的公式,執行緒數設定為 9個,但在我們實踐程序中發現如果增大執行緒數量,會顯著提高訊息的處理能力,說明 2n + 1 對于業務場景來說,并不太合適,
- 如果套用 執行緒數 = CPU核心數/(1-阻塞系數) 阻塞系數取 0.8 ,執行緒數為20
- 如果我們發現資料庫的操作耗時比較多,此時可以繼續提高阻塞系數,從而增大執行緒數量,
那我們怎么判斷需要增加更多執行緒呢?
- 其實可以用jstack命令查看一下行程的執行緒堆疊,如果發現執行緒池中大部分執行緒都處于等待獲取任務,則說明執行緒夠用
- 如果大部分執行緒都處于運行狀態,可以繼續適當調高執行緒數量,
執行緒數規劃的公式(推薦)
《Java 并發編程實戰》介紹了一個執行緒數計算的公式:

如果希望程式跑到CPU的目標利用率,需要的執行緒數公式為:

如果我期望目標利用率為90%(多核90),那么需要的執行緒數為:

把公式變個形,還可以通過執行緒數來計算CPU利用率:

雖然公式很好,但在真實的程式中,一般很難獲得準確的等待時間和計算時間,因為程式很復雜,不只是“計算” ,一段代碼中會有很多的記憶體讀寫,計算,I/O 等復合操作,精確的獲取這兩個指標很難,所以光靠公式計算執行緒數過于理想化,
真實程式中的執行緒數是沒有固定答案,先設定預期,比如我期望的CPU利用率在多少,負載在多少,GC頻率多少之類的指標后,再通過測驗不斷的調整到一個合理的執行緒數
獲取CPU核心數
- Java 獲取CPU核心數
Runtime.getRuntime().availableProcessors()//獲取邏輯核心數,如6核心12執行緒,那么回傳的是12
復制代碼
- Linux 獲取CPU核心數
# 總核數 = 物理CPU個數 X 每顆物理CPU的核數
# 總邏輯CPU數 = 物理CPU個數 X 每顆物理CPU的核數 X 超執行緒數
# 查看物理CPU個數
cat /proc/cpuinfo| grep "physical id"| sort| uniq| wc -l
# 查看每個物理CPU中core的個數(即核數)
cat /proc/cpuinfo| grep "cpu cores"| uniq
# 查看邏輯CPU的個數
cat /proc/cpuinfo| grep "processor"| wc -l
轉載請註明出處,本文鏈接:https://www.uj5u.com/qita/382812.html
標籤:其他
