Java線程池基礎(chǔ)教程_第1頁
Java線程池基礎(chǔ)教程_第2頁
Java線程池基礎(chǔ)教程_第3頁
Java線程池基礎(chǔ)教程_第4頁
Java線程池基礎(chǔ)教程_第5頁
已閱讀5頁,還剩24頁未讀 繼續(xù)免費閱讀

下載本文檔

版權(quán)說明:本文檔由用戶提供并上傳,收益歸屬內(nèi)容提供方,若內(nèi)容存在侵權(quán),請進行舉報或認領(lǐng)

文檔簡介

線程池詳解(ThreadPoolExecutor)前言在實現(xiàn)異步時,基本都是使用線程池來實現(xiàn),線程池在工作應(yīng)用的還是比較頻繁的,本文將就線程池的使用、相關(guān)原理和主要方法源碼進行深入講解學習。線程池的基本使用importjava.util.ArrayList;importjava.util.List;importjava.util.concurrent.Callable;importjava.util.concurrent.ExecutorService;importjava.util.concurrent.Executors;importjava.util.concurrent.Future;importjava.util.concurrent.FutureTask;importjava.util.concurrent.LinkedBlockingQueue;importjava.util.concurrent.ScheduledExecutorService;importjava.util.concurrent.ThreadPoolExecutor;importjava.util.concurrent.TimeUnit;publicclassThreadPoolExecutorTest{/***創(chuàng)建一個線程池(完整入?yún)?:*核心線程數(shù)為5(corePoolSize),*最大線程數(shù)為10(maximumPoolSize),*存活時間為60分鐘(keepAliveTime),*工作隊列為LinkedBlockingQueue(workQueue),*線程工廠為默認的DefaultThreadFactory(threadFactory),*飽和策略(拒絕策略)為AbortPolicy:拋出異常(handler).*/privatestaticExecutorServiceTHREAD_POOL=newThreadPoolExecutor(5,10,60,TimeUnit.MINUTES,newLinkedBlockingQueue<Runnable>(),Executors.defaultThreadFactory(),newThreadPoolExecutor.AbortPolicy());/***只有一個線程的線程池沒有超時時間,工作隊列使用無界的LinkedBlockingQueue*/privatestaticExecutorServicesingleThreadExecutor=Executors.newSingleThreadExecutor();//privatestaticExecutorServicesingleThreadExecutor=Executors.newSingleThreadExecutor(Executors.defaultThreadFactory());/***有固定線程的線程池(即corePoolSize=maximumPoolSize)沒有超時時間,*工作隊列使用無界的LinkedBlockingQueue*/privatestaticExecutorServicefixedThreadPool=Executors.newFixedThreadPool(5);//privatestaticExecutorServicefixedThreadPool=Executors.newFixedThreadPool(5,Executors.defaultThreadFactory());/***大小不限的線程池核心線程數(shù)為0,最大線程數(shù)為Integer.MAX_VALUE,存活時間為60秒該線程池可以無限擴展,*并且當需求降低時會自動收縮,工作隊列使用同步移交SynchronousQueue.*/privatestaticExecutorServicecachedThreadPool=Executors.newCachedThreadPool();//privatestaticExecutorServicecachedThreadPool=Executors.newCachedThreadPool(Executors.defaultThreadFactory());/***給定的延遲之后運行任務(wù),或者定期執(zhí)行任務(wù)的線程池*/privatestaticScheduledExecutorServicescheduledThreadPool=Executors.newScheduledThreadPool(5);//privatestaticScheduledExecutorServicescheduledThreadPool=Executors.newScheduledThreadPool(5,Executors.defaultThreadFactory());publicstaticvoidmain(Stringargs[])throwsException{/***例子1:沒有返回結(jié)果的異步任務(wù)*/THREAD_POOL.submit(newRunnable(){@Overridepublicvoidrun(){//dosomethingSystem.out.println("沒有返回結(jié)果的異步任務(wù)");}});/***例子2:有返回結(jié)果的異步任務(wù)*/Future<List<String>>future=THREAD_POOL.submit(newCallable<List<String>>(){@OverridepublicList<String>call(){List<String>result=newArrayList<>();result.add("JoonWhee");returnresult;}});List<String>result=future.get();//獲取返回結(jié)果System.out.println("有返回結(jié)果的異步任務(wù):"+result);/***例子3:*有延遲的,周期性執(zhí)行異步任務(wù)*本例子為:延遲1秒,每2秒執(zhí)行1次*/scheduledThreadPool.scheduleAtFixedRate(newRunnable(){@Overridepublicvoidrun(){System.out.println("thisis"+Thread.currentThread().getName());}},1,2,TimeUnit.SECONDS);/***例子4:FutureTask的使用*/Callable<String>task=newCallable<String>(){publicStringcall(){return"JoonWhee";}};FutureTask<String>futureTo=newFutureTask<String>(task);THREAD_POOL.submit(futureTo);System.out.println(futureTo.get());//獲取返回結(jié)果//System.out.println(futureTo.get(3,TimeUnit.SECONDS));//超時時間為3秒}}線程池的定義和優(yōu)點線程池,從字面含義來看,是指管理一組同構(gòu)工作線程的資源池。線程池是與工作隊列密切相關(guān)的,其中在工作隊列中保存了所有等待執(zhí)行的任務(wù)。工作者線程的任務(wù)很簡單:從工作隊列中獲取一個任務(wù),執(zhí)行任務(wù),然后返回線程池并等待下一個任務(wù)。“在線程池中執(zhí)行任務(wù)“”比“為每個線程分配一個任務(wù)”優(yōu)勢更多。通過重用現(xiàn)有的線程而不是創(chuàng)建線程,可以在處理多個請求時分攤在線程創(chuàng)建和銷毀過程中產(chǎn)生的巨大開銷。另一個額外的好處是,當請求到達時,工作線程通常已經(jīng)存在,因此不會由于等待創(chuàng)建線程而延遲任務(wù)的執(zhí)行,從而提高了響應(yīng)性。通過適當?shù)恼{(diào)整線程池的大小,可以創(chuàng)建足夠的線程以便使處理器保持忙碌狀態(tài),同時還可以防止過多線程相互競爭資源而使應(yīng)用程序耗盡內(nèi)存或失敗。線程池的工作流程默認情況下,創(chuàng)建完線程池后并不會立即創(chuàng)建線程,而是等到有任務(wù)提交時才會創(chuàng)建線程來進行處理。(除非調(diào)用prestartCoreThread或prestartAllCoreThreads方法)

當線程數(shù)小于核心線程數(shù)時,每提交一個任務(wù)就創(chuàng)建一個線程來執(zhí)行,即使當前有線程處于空閑狀態(tài),直到當前線程數(shù)達到核心線程數(shù)。

當前線程數(shù)達到核心線程數(shù)時,如果這個時候還提交任務(wù),這些任務(wù)會被放到隊列里,等到線程處理完了手頭的任務(wù)后,會來隊列中取任務(wù)處理。

當前線程數(shù)達到核心線程數(shù)并且隊列也滿了,如果這個時候還提交任務(wù),則會繼續(xù)創(chuàng)建線程來處理,直到線程數(shù)達到最大線程數(shù)。當前線程數(shù)達到最大線程數(shù)并且隊列也滿了,如果這個時候還提交任務(wù),則會觸發(fā)飽和策略。

如果某個線程的控線時間超過了keepAliveTime,那么將被標記為可回收的,并且當前線程池的當前大小超過了核心線程數(shù)時,這個線程將被終止。

工作隊列如果新請求的到達速率超過了線程池的處理速率,那么新到來的請求將累積起來。在線程池中,這些請求會在一個由Executor管理的Runnable隊列中等待,而不會像線程那樣去競爭CPU資源。常見的工作隊列有以下幾種,前三種用的最多。ArrayBlockingQueue:列表形式的工作隊列,必須要有初始隊列大小,有界隊列,先進先出。LinkedBlockingQueue:鏈表形式的工作隊列,可以選擇設(shè)置初始隊列大小,有界/無界隊列,先進先出。SynchronousQueue:SynchronousQueue不是一個真正的隊列,而是一種在線程之間移交的機制。要將一個元素放入SynchronousQueue中,必須有另一個線程正在等待接受這個元素.如果沒有線程等待,并且線程池的當前大小小于最大值,那么ThreadPoolExecutor將創(chuàng)建一個線程,否則根據(jù)飽和策略,這個任務(wù)將被拒絕。使用直接移交將更高效,因為任務(wù)會直接移交給執(zhí)行它的線程,而不是被首先放在隊列中,然后由工作者線程從隊列中提取任務(wù).只有當線程池是無解的或者可以拒絕任務(wù)時,SynchronousQueue才有實際價值.PriorityBlockingQueue:優(yōu)先級隊列,有界隊列,根據(jù)優(yōu)先級來安排任務(wù),任務(wù)的優(yōu)先級是通過自然順序或Comparator(如果任務(wù)實現(xiàn)了Comparator)來定義的。

DelayedWorkQueue:延遲的工作隊列,無界隊列。

飽和策略(拒絕策略)當有界隊列被填滿后,飽和策略開始發(fā)揮作用。ThreadPoolExecutor的飽和策略可以通過調(diào)用setRejectedExecutionHandler來修改。(如果某個任務(wù)被提交到一個已被關(guān)閉的Executor時,也會用到飽和策略)。飽和策略有以下四種,一般使用默認的AbortPolicy。AbortPolicy:中止策略。默認的飽和策略,拋出未檢查的RejectedExecutionException。調(diào)用者可以捕獲這個異常,然后根據(jù)需求編寫自己的處理代碼。

DiscardPolicy:拋棄策略。當新提交的任務(wù)無法保存到隊列中等待執(zhí)行時,該策略會悄悄拋棄該任務(wù)。

DiscardOldestPolicy:拋棄最舊的策略。當新提交的任務(wù)無法保存到隊列中等待執(zhí)行時,則會拋棄下一個將被執(zhí)行的任務(wù),然后嘗試重新提交新的任務(wù)。(如果工作隊列是一個優(yōu)先隊列,那么“拋棄最舊的”策略將導致拋棄優(yōu)先級最高的任務(wù),因此最好不要將“拋棄最舊的”策略和優(yōu)先級隊列放在一起使用)。

CallerRunsPolicy:調(diào)用者運行策略。該策略實現(xiàn)了一種調(diào)節(jié)機制,該策略既不會拋棄任務(wù),也不會拋出異常,而是將某些任務(wù)回退到調(diào)用者(調(diào)用線程池執(zhí)行任務(wù)的主線程),從而降低新任務(wù)的流程。它不會在線程池的某個線程中執(zhí)行新提交的任務(wù),而是在一個調(diào)用了execute的線程中執(zhí)行該任務(wù)。當線程池的所有線程都被占用,并且工作隊列被填滿后,下一個任務(wù)會在調(diào)用execute時在主線程中執(zhí)行(調(diào)用線程池執(zhí)行任務(wù)的主線程)。由于執(zhí)行任務(wù)需要一定時間,因此主線程至少在一段時間內(nèi)不能提交任務(wù),從而使得工作者線程有時間來處理完正在執(zhí)行的任務(wù)。在這期間,主線程不會調(diào)用accept,因此到達的請求將被保存在TCP層的隊列中。如果持續(xù)過載,那么TCP層將最終發(fā)現(xiàn)它的請求隊列被填滿,因此同樣會開始拋棄請求。當服務(wù)器過載后,這種過載情況會逐漸向外蔓延開來——從線程池到工作隊列到應(yīng)用程序再到TCP層,最終達到客戶端,導致服務(wù)器在高負載下實現(xiàn)一種平緩的性能降低。

線程工廠每當線程池需要創(chuàng)建一個線程時,都是通過線程工廠方法來完成的。在ThreadFactory中只定義了一個方法newThread,每當線程池需要創(chuàng)建一個新線程時都會調(diào)用這個方法。Executors提供的線程工廠有兩種,一般使用默認的,當然如果有特殊需求,也可以自己定制。DefaultThreadFactory:默認線程工廠,創(chuàng)建一個新的、非守護的線程,并且不包含特殊的配置信息。PrivilegedThreadFactory:通過這種方式創(chuàng)建出來的線程,將與創(chuàng)建privilegedThreadFactory的線程擁有相同的訪問權(quán)限、AccessControlContext、ContextClassLoader。如果不使用privilegedThreadFactory,線程池創(chuàng)建的線程將從在需要新線程時調(diào)用execute或submit的客戶程序中繼承訪問權(quán)限。自定義線程工廠:可以自己實現(xiàn)ThreadFactory接口來定制自己的線程工廠方法。

ThreadPoolExecutor源碼解析幾個點了解這幾個點,有助于你閱讀下面的源碼解釋。下面的源碼解讀中提到的運行狀態(tài)就是runState,有效的線程數(shù)就是workerCount,內(nèi)容比較多,所以可能兩種寫法都用到。運行狀態(tài)的一些定義:RUNNING:接受新任務(wù)并處理排隊任務(wù);SHUTDOWN:不接受新任務(wù),但處理排隊任務(wù);STOP:不接受新任務(wù),不處理排隊任務(wù),并中斷正在進行的任務(wù);TIDYING:所有任務(wù)已經(jīng)終止,workerCount為零,線程轉(zhuǎn)換到狀態(tài)TIDYING將運行terminate()鉤子方法;TERMINATED:terminated()已經(jīng)完成,該方法執(zhí)行完畢代表線程池已經(jīng)完全終止。運行狀態(tài)之間并不是隨意轉(zhuǎn)換的,大多數(shù)狀態(tài)都只能由固定的狀態(tài)轉(zhuǎn)換而來,轉(zhuǎn)換關(guān)系見第4點~第8點。RUNNING->SHUTDOWN:在調(diào)用shutdown()時,可能隱含在finalize()。(RUNNINGorSHUTDOWN)->STOP:調(diào)用shutdownNow()。SHUTDOWN->TIDYING:當隊列和線程池都是空的時。

STOP->TIDYING:當線程池為空時。

TIDYING->TERMINATED:當terminate()方法完成時。

基礎(chǔ)屬性(很重要)/***主池控制狀態(tài)ctl是包含兩個概念字段的原子整數(shù):workerCount:指有效的線程數(shù)量;*runState:指運行狀態(tài),運行,關(guān)閉等。為了將workerCount和runState用1個int來表示,*我們限制workerCount范圍為(2^29)-1,即用int的低29位用來表示workerCount,*用int的高3位用來表示runState,這樣workerCount和runState剛好用int可以完整表示。*///初始化時有效的線程數(shù)為0,此時ctl為:10100000000000000000000000000000privatefinalAtomicIntegerctl=newAtomicInteger(ctlOf(RUNNING,0));//高3位用來表示運行狀態(tài),此值用于運行狀態(tài)向左移動的位數(shù),即29位privatestaticfinalintCOUNT_BITS=Integer.SIZE-3;//線程數(shù)容量,低29位表示有效的線程數(shù),00011111111111111111111111111111privatestaticfinalintCAPACITY=(1<<COUNT_BITS)-1;/***大小關(guān)系:RUNNING<SHUTDOWN<STOP<TIDYING<TERMINATED,*源碼中頻繁使用大小關(guān)系來作為條件判斷。*10100000000000000000000000000000運行*01100000000000000000000000000000關(guān)閉*01100000000000000000000000000000停止*01100000000000000000000000000000整理*01100000000000000000000000000000終止*/privatestaticfinalintRUNNING=-1<<COUNT_BITS;//運行privatestaticfinalintSHUTDOWN=0<<COUNT_BITS;//關(guān)閉privatestaticfinalintSTOP=1<<COUNT_BITS;//停止privatestaticfinalintTIDYING=2<<COUNT_BITS;//整理privatestaticfinalintTERMINATED=3<<COUNT_BITS;//終止/***得到運行狀態(tài):入?yún)為ctl的值,~CAPACITY高3位為1低29位全為0,*因此運算結(jié)果為ctl的高3位,也就是運行狀態(tài)*/privatestaticintrunStateOf(intc){returnc&~CAPACITY;}/***得到有效的線程數(shù):入?yún)為ctl的值,CAPACITY高3為為0,*低29位全為1,因此運算結(jié)果為ctl的低29位,也就是有效的線程數(shù)*/privatestaticintworkerCountOf(intc){returnc&CAPACITY;}/***得到ctl的值:高3位的運行狀態(tài)和低29位的有效線程數(shù)進行或運算,*組合成一個完成的32位數(shù)*/privatestaticintctlOf(intrs,intwc){returnrs|wc;}//狀態(tài)c是否小于sprivatestaticbooleanrunStateLessThan(intc,ints){returnc<s;}//狀態(tài)c是否大于等于sprivatestaticbooleanrunStateAtLeast(intc,ints){returnc>=s;}//狀態(tài)c是否為RUNNING(小于SHUTDOWN的狀態(tài)只有RUNNING)privatestaticbooleanisRunning(intc){returnc<SHUTDOWN;}//使用CAS增加一個有效的線程privatebooleancompareAndIncrementWorkerCount(intexpect){returnpareAndSet(expect,expect+1);}//使用CAS減少一個有效的線程privatebooleancompareAndDecrementWorkerCount(intexpect){returnpareAndSet(expect,expect-1);}//減少一個有效的線程privatevoiddecrementWorkerCount(){do{}while(!compareAndDecrementWorkerCount(ctl.get()));}//工作隊列privatefinalBlockingQueue<Runnable>workQueue;//鎖privatefinalReentrantLockmainLock=newReentrantLock();//包含線程池中的所有工作線程,只有在mainLock的情況下才能訪問,Worker集合privatefinalHashSet<Worker>workers=newHashSet<Worker>();privatefinalConditiontermination=mainLock.newCondition();//跟蹤線程池的最大到達大小,僅在mainLock下訪問privateintlargestPoolSize;//總的完成的任務(wù)數(shù)privatelongcompletedTaskCount;//線程工廠,用于創(chuàng)建線程privatevolatileThreadFactorythreadFactory;//拒絕策略privatevolatileRejectedExecutionHandlerhandler;/***線程超時時間,當線程數(shù)超過corePoolSize時生效,*如果有線程空閑時間超過keepAliveTime,則會被終止*/privatevolatilelongkeepAliveTime;//是否允許核心線程超時,默認false,false情況下核心線程會一直存活。privatevolatilebooleanallowCoreThreadTimeOut;//核心線程數(shù)privatevolatileintcorePoolSize;//最大線程數(shù)privatevolatileintmaximumPoolSize;//默認飽和策略(拒絕策略),拋異常privatestaticfinalRejectedExecutionHandlerdefaultHandler=newAbortPolicy();privatestaticfinalRuntimePermissionshutdownPerm=newRuntimePermission("modifyThread");/***Worker類,每個Worker包含一個線程、一個初始任務(wù)、一個任務(wù)計算器*/privatefinalclassWorkerextendsAbstractQueuedSynchronizerimplementsRunnable{privatestaticfinallongserialVersionUID=6138294804551838833L;finalThreadthread;//Worker對應(yīng)的線程RunnablefirstTask;//運行的初始任務(wù)。volatilelongcompletedTasks;//每個線程的任務(wù)計數(shù)器Worker(RunnablefirstTask){setState(-1);//禁止中斷,直到runWorkerthis.firstTask=firstTask;//設(shè)置為初始任務(wù)//使用當前線程池的線程工廠創(chuàng)建一個線程this.thread=getThreadFactory().newThread(this);}//將主運行循環(huán)委托給外部runWorkerpublicvoidrun(){runWorker(this);}//Lockmethods////Thevalue0representstheunlockedstate.//Thevalue1representsthelockedstate./***通過AQS的同步狀態(tài)來實現(xiàn)鎖機制。state為0時代表鎖未被獲取(解鎖狀態(tài)),*state為1時代表鎖已經(jīng)被獲取(加鎖狀態(tài))。*/protectedbooleanisHeldExclusively(){//returngetState()!=0;}protectedbooleantryAcquire(intunused){//嘗試獲取鎖if(compareAndSetState(0,1)){//使用CAS嘗試將state設(shè)置為1,即嘗試獲取鎖//成功將state設(shè)置為1,則當前線程擁有獨占訪問權(quán)setExclusiveOwnerThread(Thread.currentThread());returntrue;}returnfalse;}protectedbooleantryRelease(intunused){//嘗試釋放鎖setExclusiveOwnerThread(null);//釋放獨占訪問權(quán):即將獨占訪問線程設(shè)為nullsetState(0);//解鎖:將state設(shè)置為0returntrue;}publicvoidlock(){acquire(1);}//加鎖publicbooleantryLock(){returntryAcquire(1);}//嘗試加鎖publicvoidunlock(){release(1);}//解鎖publicbooleanisLocked(){returnisHeldExclusively();}//是否為加鎖狀態(tài)voidinterruptIfStarted(){//如果線程啟動了,則進行中斷Threadt;if(getState()>=0&&(t=thread)!=null&&!t.isInterrupted()){try{errupt();}catch(SecurityExceptionignore){}}}}execute方法使用線程池的submit方法提交任務(wù)時,會走到該方法,該方法也是線程池最重要的方法。

publicvoidexecute(Runnablecommand){if(command==null)//為空校驗thrownewNullPointerException();intc=ctl.get();//拿到當前的ctl值if(workerCountOf(c)<corePoolSize){//如果有效的線程數(shù)小于核心線程數(shù)if(addWorker(command,true))//則新建一個線程來處理任務(wù)(核心線程)return;c=ctl.get();//拿到當前的ctl值}//走到這里說明有效的線程數(shù)已經(jīng)>=核心線程數(shù)if(isRunning(c)&&workQueue.offer(command)){//如果當前狀態(tài)是運行,嘗試將任務(wù)放入工作隊列intrecheck=ctl.get();//再次拿到當前的ctl值//如果再次檢查狀態(tài)不是運行,則將剛才添加到工作隊列的任務(wù)移除if(!isRunning(recheck)&&remove(command))reject(command);//并調(diào)用拒絕策略elseif(workerCountOf(recheck)==0)//如果再次檢查時,有效的線程數(shù)為0,addWorker(null,false);//則新建一個線程(非核心線程)}//走到這里說明工作隊列已滿elseif(!addWorker(command,false))//嘗試新建一個線程來處理任務(wù)(非核心)reject(command);//如果失敗則調(diào)用拒絕策略}該方法就是對應(yīng)上文的線程池的工作流程。主要調(diào)用到的方法為addWorker(見下文addWorker方法解讀)。

addWorker方法方法主要目的就是使用入?yún)⒅械膄irstTask和當前線程添加一個Worker,前面的for循環(huán)主要是對當前線程池的運行狀態(tài)和有效的線程數(shù)進行一些校驗,校驗邏輯比較繞,可以參考注釋進行理解。該方法涉及到的其他方法有addWorkerFailed(見下文addWorkerFailed源碼解讀);還有就是Worker的線程啟動時,會調(diào)用Worker里的run方法,執(zhí)行runWorker(this)方法(見下文runWorker源碼解讀)。

addWorkerFailed方法/***Rollsbacktheworkerthreadcreation.*-removesworkerfromworkers,ifpresent*-decrementsworkercount*-rechecksfortermination,incasetheexistenceofthis*workerwasholdinguptermination*/privatevoidaddWorkerFailed(Workerw){//回滾Worker的添加,就是將Worker移除finalReentrantLockmainLock=this.mainLock;mainLock.lock();try{if(w!=null)workers.remove(w);//移除WorkerdecrementWorkerCount();//有效線程數(shù)-1tryTerminate();//有worker線程移除,可能是最后一個線程退出需要嘗試終止線程池}finally{mainLock.unlock();}}該方法很簡單,就是移除入?yún)⒅械腤orker并將workerCount-1,最后調(diào)用tryTerminate嘗試終止線程池,tryTerminate見下文對應(yīng)方法源碼解讀。runWorker方法上文addWork方法里說道,當Worker里的線程啟動時,就會調(diào)用該方法。/***Worker的線程開始執(zhí)行任務(wù)*/finalvoidrunWorker(Workerw){Threadwt=Thread.currentThread();//獲取當前線程Runnabletask=w.firstTask;//拿到Worker的初始任務(wù)w.firstTask=null;w.unlock();//allowinterruptsbooleancompletedAbruptly=true;//Worker是不是因異常而死亡try{while(task!=null||(task=getTask())!=null){//Worker取任務(wù)執(zhí)行w.lock();//加鎖/**如果線程池停止,確保線程中斷;如果不是,確保線程不被中斷。*在第二種情況下進行重新檢查,以便在清除中斷的同時處理shutdownNow競爭*線程池停止指運行狀態(tài)為STOP/TIDYING/TERMINATED中的一種*/if((runStateAtLeast(ctl.get(),STOP)||//判斷線程池運行狀態(tài)(Terrupted()&&//重新檢查runStateAtLeast(ctl.get(),STOP)))&&//再次判斷線程池運行狀態(tài)!wt.isInterrupted())//走到這里代表線程池運行狀態(tài)為停止,檢查wt是否中斷errupt();//線程池的狀態(tài)為停止并且wt不為中斷,則將wt中斷try{beforeExecute(wt,task);//執(zhí)行beforeExecute(默認空,需要自己重寫)Throwablethrown=null;try{task.run();//執(zhí)行任務(wù)}catch(RuntimeExceptionx){thrown=x;throwx;//如果拋異常,則completedAbruptly為true}catch(Errorx){thrown=x;throwx;}catch(Throwablex){thrown=x;thrownewError(x);}finally{afterExecute(task,thrown);//執(zhí)行afterExecute(需要自己重寫)}}finally{task=null;//將執(zhí)行完的任務(wù)清空pletedTasks++;//Worker完成任務(wù)數(shù)+1w.unlock();}}completedAbruptly=false;//如果執(zhí)行到這里,則worker是正常退出}finally{processWorkerExit(w,completedAbruptly);//調(diào)用processWorkerExit方法}}該方法為Worker線程開始執(zhí)行任務(wù),首先執(zhí)行當初創(chuàng)建Worker時的初始任務(wù),接著從工作隊列中獲取任務(wù)執(zhí)行。主要涉及兩個方法:獲取任務(wù)的方法getTask(見下文getTask源碼解讀)和執(zhí)行Worker退出的方法processWorkerExit(見下文processWorkerExit源碼解讀)。注:processWorkerExit在處理正常Worker退出時,沒有對workerCount-1,而是在getTask方法中進行workerCount-1。

getTask方法privateRunnablegetTask(){//Worker從工作隊列獲取任務(wù)booleantimedOut=false;//poll方法取任務(wù)是否超時for(;;){//無線循環(huán)intc=ctl.get();//ctlintrs=runStateOf(c);//當前運行狀態(tài)//如果線程池運行狀態(tài)為停止,或者可以停止(狀態(tài)為SHUTDOWN并且隊列為空)//則返回null,代表當前Worker需要移除if(rs>=SHUTDOWN&&(rs>=STOP||workQueue.isEmpty())){decrementWorkerCount();//將workerCount-1//返回null前將workerCount-1,//因此processWorkerExit中completedAbruptly=false時無需再減returnnull;}intwc=workerCountOf(c);//當前的workerCount//判斷當前Worker是否可以被移除,即當前Worker是否可以一直等待任務(wù)。//如果allowCoreThreadTimeOut為true,或者workerCount大于核心線程數(shù),//則當前線程是有超時時間的(keepAliveTime),無法一直等待任務(wù)。booleantimed=allowCoreThreadTimeOut||wc>corePoolSize;//如果wc超過最大線程數(shù)或者當前線程會超時并且已經(jīng)超時,//并且wc>1或者工作隊列為空,則返回null,代表當前Worker需要移除if((wc>maximumPoolSize||(timed&&timedOut))&&(wc>1||workQueue.isEmpty())){//確保有Worker可以移除if(compareAndDecrementWorkerCount(c))//返回null前將workerCount-1,//因此processWorkerExit中completedAbruptly=false時無需再減returnnull;continue;}try{//根據(jù)線程是否會超時調(diào)用相應(yīng)的方法,poll為帶超時的獲取任務(wù)方法//take()為不帶超時的獲取任務(wù)方法,會一直阻塞直到獲取到任務(wù)Runnabler=timed?workQueue.poll(keepAliveTime,TimeUnit.NANOSECONDS):workQueue.take();if(r!=null)returnr;timedOut=true;//走到這代表當前線程獲取任務(wù)超時}catch(InterruptedExceptionretry){timedOut=false;//被中斷}}}Worker從工作隊列獲取任務(wù),如果allowCoreThreadTimeOut為false并且

workerCount<=corePoolSize,則這些核心線程永遠存活,并且一直在嘗試獲取工作隊列的任務(wù);否則,線程會有超時時間(keepAliveTime),當在keepAliveTime時間內(nèi)獲取不到任務(wù),該線程的Worker會被移除。

Worker移除的過程:getTask方法返回null,導致runWorker方法中跳出while循環(huán),調(diào)用processWorkerExit方法將Worker移除。注意:在返回null的之前,已經(jīng)將workerCount-1,因此在processWorkerExit中,completedAbruptly=false的情況(即正常超時退出)不需要再將workerCount-1。

processWorkerExit方法privatevoidprocessWorkerExit(Workerw,booleancompletedAbruptly){//Worker的退出//如果Worker是異常死亡(completedAbruptly=true),則workerCount-1;//如果completedAbruptly為false的時候(正常超時退出),則代表task=getTask()等于null,//getTask()方法中返回null的地方,都已經(jīng)將workerCount-1,所以此處無需再-1if(completedAbruptly)decrementWorkerCount();finalReentrantLockmainLock=this.mainLock;mainLock.lock();//加鎖try{completedTaskCount+=pletedTasks;//該Worker完成的任務(wù)數(shù)加到總完成的任務(wù)數(shù)workers.remove(w);//移除該Worker}finally{mainLock.unlock();}tryTerminate();//有Worker線程移除,可能是最后一個線程退出,需要嘗試終止線程池intc=ctl.get();//獲取當前的ctlif(runStateLessThan(c,STOP)){//如果線程池的運行狀態(tài)還沒停止(RUNNING或SHUTDOWN)if(!completedAbruptly){//如果Worker不是異常死亡//min為線程池的理論最小線程數(shù):如果允許核心線程超時則min為0,否則min為核心線程數(shù)intmin=allowCoreThreadTimeOut?0:corePoolSize;//如果min為0,工作隊列不為空,將min設(shè)置為1,確保至少有1個Worker來處理隊列里的任務(wù)if(min==0&&!w

溫馨提示

  • 1. 本站所有資源如無特殊說明,都需要本地電腦安裝OFFICE2007和PDF閱讀器。圖紙軟件為CAD,CAXA,PROE,UG,SolidWorks等.壓縮文件請下載最新的WinRAR軟件解壓。
  • 2. 本站的文檔不包含任何第三方提供的附件圖紙等,如果需要附件,請聯(lián)系上傳者。文件的所有權(quán)益歸上傳用戶所有。
  • 3. 本站RAR壓縮包中若帶圖紙,網(wǎng)頁內(nèi)容里面會有圖紙預覽,若沒有圖紙預覽就沒有圖紙。
  • 4. 未經(jīng)權(quán)益所有人同意不得將文件中的內(nèi)容挪作商業(yè)或盈利用途。
  • 5. 人人文庫網(wǎng)僅提供信息存儲空間,僅對用戶上傳內(nèi)容的表現(xiàn)方式做保護處理,對用戶上傳分享的文檔內(nèi)容本身不做任何修改或編輯,并不能對任何下載內(nèi)容負責。
  • 6. 下載文件中如有侵權(quán)或不適當內(nèi)容,請與我們聯(lián)系,我們立即糾正。
  • 7. 本站不保證下載資源的準確性、安全性和完整性, 同時也不承擔用戶因使用這些下載資源對自己和他人造成任何形式的傷害或損失。

評論

0/150

提交評論