最新国产好看的视频,伊人天堂AV在线,国产Aaaaaa视频,蜜臀视频在线观看一区,人妻av色图,密臀久久久精品影片,青青视频免费观看毛片,久草在线观看视,国产三级精品色情在线

支持生產(chǎn)阻塞的Java線程池

 更新時間:2014年04月17日 09:11:44   作者:  
在各種并發(fā)編程模型中,生產(chǎn)者-消費者模式大概是最常用的了。在實際工作中,對于生產(chǎn)消費的速度,通常需要做一下權衡

通常來說,生產(chǎn)任務的速度要大于消費的速度。一個細節(jié)問題是,隊列長度,以及如何匹配生產(chǎn)和消費的速度。

一個典型的生產(chǎn)者-消費者模型如下:

 

在并發(fā)環(huán)境下利用J.U.C提供的Queue實現(xiàn)可以很方便地保證生產(chǎn)和消費過程中的線程安全。這里需要注意的是,Queue必須設置初始容量,防止生產(chǎn)者生產(chǎn)過快導致隊列長度暴漲,最終觸發(fā)OutOfMemory。

對于一般的生產(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和Consumer部分已經(jīng)幫我們實現(xiàn)好了,并且直接采用線程池的實現(xiàn)還有很多優(yōu)勢,例如線程數(shù)的動態(tài)調(diào)整等。

但問題在于,即便你在構造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)者阻塞,那就顯得復雜了。

相關文章

最新評論

武汉市| 固原市| 凤山县| 宁明县| 白河县| 乌什县| 乌鲁木齐市| 拜城县| 峡江县| 西峡县| 淮滨县| 巴青县| 宜城市| 柳林县| 拉孜县| 彭泽县| 山东| 宁明县| 凤台县| 长海县| 台湾省| 崇州市| 宁武县| 中山市| 大港区| 潢川县| 甘谷县| 富裕县| 措勤县| 平乡县| 海阳市| 德保县| 锡林浩特市| 阿拉善左旗| 巴南区| 双辽市| 合江县| 永靖县| 游戏| 明水县| 莱州市|