支持生產(chǎn)阻塞的Java線程池
通常來說,生產(chǎn)任務的速度要大于消費的速度。一個細節(jié)問題是,隊列長度,以及如何匹配生產(chǎn)和消費的速度。
一個典型的生產(chǎn)者-消費者模型如下:
![]() |
對于一般的生產(chǎn)快于消費的情況。當隊列已滿時,我們并不希望有任何任務被忽略或得不到執(zhí)行,此時生產(chǎn)者可以等待片刻再提交任務,更好的做法是,把生產(chǎn)者阻塞在提交任務的方法上,待隊列未滿時繼續(xù)提交任務,這樣就沒有浪費的空轉(zhuǎn)時間了。阻塞這一點也很容易,BlockingQueue就是為此打造的,ArrayBlockingQueue和LinkedBlockingQueue在構造時都可以提供容量做限制,其中LinkedBlockingQueue是在實際操作隊列時在每次拿到鎖以后判斷容量。
更進一步,當隊列為空時,消費者拿不到任務,可以等一會兒再拿,更好的做法是,用BlockingQueue的take方法,阻塞等待,當有任務時便可以立即獲得執(zhí)行,建議調(diào)用take的帶超時參數(shù)的重載方法,超時后線程退出。這樣當生產(chǎn)者事實上已經(jīng)停止生產(chǎn)時,不至于讓消費者無限等待。
于是一個高效的支持阻塞的生產(chǎn)消費模型就實現(xiàn)了。
等一下,既然J.U.C已經(jīng)幫我們實現(xiàn)了線程池,為什么還要采用這一套東西?直接用ExecutorService不是更方便?
我們來看一下ThreadPoolExecutor的基本結構:
![]() |
但問題在于,即便你在構造ThreadPoolExecutor時手動指定了一個BlockingQueue作為隊列實現(xiàn),事實上當隊列滿時,execute方法并不會阻塞,原因在于ThreadPoolExecutor調(diào)用的是BlockingQueue非阻塞的offer方法:
public void execute(Runnable command) {
if (command == null)
throw new NullPointerException();
if (poolSize >= corePoolSize || !addIfUnderCorePoolSize(command)) {
if (runState == RUNNING && workQueue.offer(command)) {
if (runState != RUNNING || poolSize == 0)
ensureQueuedTaskHandled(command);
}
else if (!addIfUnderMaximumPoolSize(command))
reject(command); // is shutdown or saturated
}
}
這時候就需要做一些事情來達成一個結果:當生產(chǎn)者提交任務,而隊列已滿時,能夠讓生產(chǎn)者阻塞住,等待任務被消費。
關鍵在于,在并發(fā)環(huán)境下,隊列滿不能由生產(chǎn)者去判斷,不能調(diào)用ThreadPoolExecutor.getQueue().size()來判斷隊列是否滿。
線程池的實現(xiàn)中,當隊列滿時會調(diào)用構造時傳入的RejectedExecutionHandler去拒絕任務的處理。默認的實現(xiàn)是AbortPolicy,直接拋出一個RejectedExecutionException。
幾種拒絕策略在這里就不贅述了,這里和我們的需求比較接近的是CallerRunsPolicy,這種策略會在隊列滿時,讓提交任務的線程去執(zhí)行任務,相當于讓生產(chǎn)者臨時去干了消費者干的活兒,這樣生產(chǎn)者雖然沒有被阻塞,但提交任務也會被暫停。
public static class CallerRunsPolicy implements RejectedExecutionHandler {
/**
* Creates a <tt>CallerRunsPolicy</tt>.
*/
public CallerRunsPolicy() { }
/**
* Executes task r in the caller's thread, unless the executor
* has been shut down, in which case the task is discarded.
* @param r the runnable task requested to be executed
* @param e the executor attempting to execute this task
*/
public void rejectedExecution(Runnable r, ThreadPoolExecutor e) {
if (!e.isShutdown()) {
r.run();
}
}
}
但這種策略也有隱患,當生產(chǎn)者較少時,生產(chǎn)者消費任務的時間里,消費者可能已經(jīng)把任務都消費完了,隊列處于空狀態(tài),當生產(chǎn)者執(zhí)行完任務后才能再繼續(xù)生產(chǎn)任務,這個過程中可能導致消費者線程的饑餓。
參考類似的思路,最簡單的做法,我們可以直接定義一個RejectedExecutionHandler,當隊列滿時改為調(diào)用BlockingQueue.put來實現(xiàn)生產(chǎn)者的阻塞:
new RejectedExecutionHandler() {
@Override
public void rejectedExecution(Runnable r, ThreadPoolExecutor executor) {
if (!executor.isShutdown()) {
try {
executor.getQueue().put(r);
} catch (InterruptedException e) {
// should not be interrupted
}
}
}
};
這樣,我們就無需再關心Queue和Consumer的邏輯,只要把精力集中在生產(chǎn)者和消費者線程的實現(xiàn)邏輯上,只管往線程池提交任務就行了。
相比最初的設計,這種方式的代碼量能減少不少,而且能避免并發(fā)環(huán)境的很多問題。當然,你也可以采用另外的手段,例如在提交時采用信號量做入口限制等,但是如果僅僅是要讓生產(chǎn)者阻塞,那就顯得復雜了。
相關文章
Spring啟動過程中實例化部分代碼的分析之Bean的推斷構造方法
這篇文章主要介紹了Spring啟動過程中實例化部分代碼的分析之Bean的推斷構造方法,實例化這一步便是在doCreateBean方法的?instanceWrapper?=?createBeanInstance(beanName,?mbd,?args);這段代碼中,本文通過實例代碼給大家介紹的非常詳細,需要的朋友參考下吧2022-09-09
Java編程之多線程死鎖與線程間通信簡單實現(xiàn)代碼
這篇文章主要介紹了Java編程之多線程死鎖與線程間通信簡單實現(xiàn)代碼,具有一定參考價值,需要的朋友可以了解下。2017-10-10
淺析Java中Map與HashMap,Hashtable,HashSet的區(qū)別
HashMap和Hashtable兩個類都實現(xiàn)了Map接口,二者保存K-V對(key-value對);HashSet則實現(xiàn)了Set接口,性質(zhì)類似于集合2013-09-09
SpringBoot使用ApplicationEvent&Listener完成業(yè)務解耦
這篇文章主要介紹了SpringBoot使用ApplicationEvent&Listener完成業(yè)務解耦示例詳解,有需要的朋友可以借鑒參考下,希望能夠有所幫助,祝大家多多進步,早日升職加薪2023-05-05
SpringBoot集成pf4j實現(xiàn)插件開發(fā)功能的代碼示例
pf4j是一個插件框架,用于實現(xiàn)插件的動態(tài)加載,支持的插件格式(zip、jar),本文給大家介紹了SpringBoot集成pf4j實現(xiàn)插件開發(fā)功能的示例,文中通過代碼示例給大家講解的非常詳細,需要的朋友可以參考下2024-07-07
java request.getParameter中文亂碼解決方法
今天跟大家分享幾個解決java Web開發(fā)中,request.getParameter()獲取URL中文參數(shù)亂碼的解決辦法,需要的朋友可以參考下2020-02-02
解析MapStruct轉(zhuǎn)換javaBean時出現(xiàn)的詭異事件
在項目中用到了MapStruct,對其可以轉(zhuǎn)換JavaBean特別好奇,今天小編給大家分享一個demo給大家講解MapStruct轉(zhuǎn)換javaBean時出現(xiàn)的詭異事件,感興趣的朋友一起看看吧2021-09-09
不規(guī)范使用ThreadLocal導致bug分析解決
這篇文章主要為大家介紹了不規(guī)范使用ThreadLocal導致bug分析解決,有需要的朋友可以借鑒參考下,希望能夠有所幫助,祝大家多多進步,早日升職加薪2023-01-01



