如何解決JDK執行緒池中不超過最大執行緒數下快速消費任務

首頁 > 科技

如何解決JDK執行緒池中不超過最大執行緒數下快速消費任務

來源:子非魚 釋出時間:2022-09-30 09:32

文章需要對執行緒池執行任務流程有一定的瞭解

記得之前我寫透過模版設計來解決執行緒池引數自定義痛點,然後寬哥在下面靈魂發問,也就是咱們這篇文章講到的重點

來來來,我給大家複製粘貼出來

如何解決JDK執行緒池中不超過最大執行緒數下即時快速消費任務,而不是在佇列中堆積

因為最近業務落地改造中需要執行緒池,又去看了一遍原始碼,防止線上埋雷,也再次回顧了這個問題

然後發現網上也有這種問題提問,雖然是不同的提問,但是核心思想是一致的,

業務是多變的,而JDK中的執行緒池消費流程卻是固定的,所以基於阻塞佇列、執行緒池擴充套件改變了原有流程

01、執行緒池引數

我們這裡講解以ThreadPoolExecutor#execute(Runnablerunnable)舉例,這裡先說下執行緒池的一些引數

本篇只是說明上述問題,不會對執行緒池做詳細講解publicThreadPoolExecutor(intcorePoolSize,

intmaximumPoolSize,

longkeepAliveTime,

TimeUnitunit,

BlockingQueueworkQueue,

ThreadFactorythreadFactory,

RejectedExecutionHandlerhandler){...}

corePoolSize

執行緒池中的核心執行緒數量,如果沒有全域性設定池內執行緒的過期時間,池內會維持此數量執行緒

maximumPoolSize

執行緒池中的最大執行緒數量,當核心執行緒都在執行任務,並且阻塞佇列中任務數量已滿,此時會建立非核心執行緒

keepAliveTime&unit

執行緒池中執行緒過期時間以及時間單位

workQueue

存放執行緒池內任務的阻塞佇列,如ArrayBlockingQueue、LinkedBlockingQueue...

threadFactory

建立執行緒池中執行緒的執行緒工廠,可以在建立執行緒時初始化優先順序、名稱、守護狀態...

handler

當執行緒池中全部執行緒都在執行,阻塞佇列也滿的時候,會將新增的任務執行拒絕策略,JDK執行緒池中實現了四種拒絕策略,預設AbortPolicy,丟擲異常

02、執行緒池任務新增流程

相信大家在網上看到過許多類似的執行緒池執行流程圖哈,這裡仍是簡要贅述下,原始碼如下:

publicvoidexecute(Runnablecommand){

...

intc=ctl.get;

if(workerCountOf(c)

if(addWorker(command,true))

return;

c=ctl.get;

}

if(isRunning(c)&&workQueue.offer(command)){

intrecheck=ctl.get;

if(!isRunning(recheck)&&remove(command))

reject(command);

elseif(workerCountOf(recheck)==0)

addWorker(null,false);

}elseif(!addWorker(command,false))

reject(command);

}

1、執行緒池提交任務首先判斷當前執行緒數是否大於核心執行緒數,否則建立核心執行緒執行任務

2、如果當前執行緒超過了核心執行緒數,判斷阻塞佇列是否已滿,否則將任務新增到佇列中

3、如果阻塞佇列已滿,判斷當前執行緒是否大於最大執行緒數,否則建立非核心執行緒執行任務

4、如果當前執行緒大於或即是最大執行緒數,執行拒絕策略

這道問題的意圖就是要將第二步就行改寫

如果當前執行緒大於核心執行緒數,不將任務放入阻塞佇列,而是建立非核心執行緒執行任務

舉例說明一下:

publicstaticvoidmain(String[]args){

ThreadPoolExecutorthreadPoolExecutor=

newThreadPoolExecutor(1,3,60,

TimeUnit.SECONDS,

newArrayBlockingQueue(10));

for(inti=0;i

threadPoolExecutor.execute(->{

System.out.println(Thread.currentThread.getName+"-執行任務");

LockSupport.park;

});

}

threadPoolExecutor.shutdown;

/**

*pool-1-thread-1執行任務

*/

}

看到這段程式碼,正常情況下只會有一個任務會被執行,其餘任務會被放置阻塞佇列中

而我們需要做的就是,發現池內執行緒大於核心執行緒數,不放入阻塞佇列,而是建立非核心執行緒進行消費任務

原生代碼實現參考Dubbo原始碼中EagerThreadPoolExecutor,確實能實現對應效果,這裡就不演示了,一起看一下Dubbo如何做的

03、Dubbo中實現的快速消費

Dubbo中涉及到的類有兩個,EagerThreadPoolExecutor、TaskQueue

這裡貼一下重點程式碼

3.1TaskQueue

publicclassTaskQueueextendsLinkedBlockingQueue{

...

//佇列中持有執行緒池的引用

privateEagerThreadPoolExecutorexecutor;

publicTaskQueue(intcapacity){

super(capacity);

}

publicvoidsetExecutor(EagerThreadPoolExecutorexec){

executor=exec;

}

@Override

publicbooleanoffer(Runnablerunnable){

...

//獲取執行緒池中執行緒數

intcurrentPoolThreadSize=executor.getPoolSize;

//如果有核心執行緒正在空閒,將任務加入阻塞佇列,由核心執行緒進行處理任務

if(executor.getSubmittedTaskCount

returnsuper.offer(runnable);

}

/**

*[重點]當前執行緒池執行緒數量小於最大執行緒數

*返回false,根據執行緒池原始碼,會建立非核心執行緒

*/

if(currentPoolThreadSize

returnfalse;

}

//如果當前執行緒池數量大於最大執行緒數,任務加入阻塞佇列

returnsuper.offer(runnable);

}

}

存在一個疑點,getSubmittedTaskCount是如何獲取提交任務數量的?

這裡就需要看一下EagerThreadPoolExecutor實現了,也比較簡單,只是重寫了執行緒池的兩個方法:afterExecute、execute

3.2EagerThreadPoolExecutor

publicclassEagerThreadPoolExecutorextendsThreadPoolExecutor{

/**

*taskcount

*/

privatefinalAtomicIntegersubmittedTaskCount=newAtomicInteger(0);

/**

*@returncurrenttaskswhichareexecuted

*/

publicintgetSubmittedTaskCount{

returnsubmittedTaskCount.get;

}

@Override

protectedvoidafterExecute(Runnabler,Throwablet){

submittedTaskCount.decrementAndGet;

}

@Override

publicvoidexecute(Runnablecommand){

if(command==null){

thrownewNullPointerException;

}

//donotincrementinmethodbeforeExecute!

submittedTaskCount.incrementAndGet;

try{

super.execute(command);

}catch(RejectedExecutionExceptionrx){

//retrytoofferthetaskintoqueue.

finalTaskQueuequeue=(TaskQueue)super.getQueue;

try{

if(!queue.retryOffer(command,0,TimeUnit.MILLISECONDS)){

submittedTaskCount.decrementAndGet;

thrownewRejectedExecutionException("Queuecapacityisfull.",rx);

}

}catch(InterruptedExceptionx){

submittedTaskCount.decrementAndGet;

thrownewRejectedExecutionException(x);

}

}catch(Throwablet){

//decreaseanyway

submittedTaskCount.decrementAndGet;

throwt;

}

}

}

EagerThreadPoolExecutor繼承了ThreadPoolExecutor,在execute上做了個性化設計

並在執行緒池內新增了一個任務數量的欄位,是一個原子類,新增任務時自增,任務異常及結束時遞減

這樣就能保證TaskQueue#offer(Runnablerunnable)做出邏輯處理

文章需要對執行緒池執行任務流程有一定的瞭解

記得之前我寫透過模版設計來解決執行緒池引數自定義痛點,然後寬哥在下面靈魂發問,也就是咱們這篇文章講到的重點

文章需要對執行緒池執行任務流程有一定的瞭解

記得之前我寫透過模版設計來解決執行緒池引數自定義痛點,然後寬哥在下面靈魂發問,也就是咱們這篇文章講到的重點

上一篇:男子三跪九拜... 下一篇:如何提高可轉...
猜你喜歡
熱門閱讀
同類推薦