Java 中的阻塞隊列從基礎(chǔ)到高級的深度解析
提到阻塞隊列,許多人腦海中會浮現(xiàn)出
BlockingQueue、ArrayBlockingQueue、LinkedBlockingQueue和SynchronousQueue。盡管這些實現(xiàn)看起來復(fù)雜,實際上阻塞隊列本身的概念相對簡單,真正挑戰(zhàn)在于內(nèi)部的 AQS(Abstract Queuing Synchronizer)。如果你對阻塞隊列感到陌生,希望下面的內(nèi)容能幫助你從全新角度理解它。
1、線程間通信
線程間通信是指多個線程對共享資源的操作和協(xié)調(diào)。在生產(chǎn)者-消費者模型中,生產(chǎn)者和消費者是不同種類的線程,他們對同一個資源(如隊列)進行操作。生產(chǎn)者負(fù)責(zé)向隊列中插入數(shù)據(jù),消費者負(fù)責(zé)從隊列中取出數(shù)據(jù)。
主要挑戰(zhàn)在于如何在資源達到上限時讓生產(chǎn)者等待,而在資源達到下限時讓消費者等待。線程間的這種相互調(diào)度,就是線程間通信。
以現(xiàn)實生活為例。消費者和生產(chǎn)者就像兩個線程,原本做著各自的事情,廠家管自己生產(chǎn),消費者管自己買,一般情況下彼此互不影響。900 240

但當(dāng)物資到達某個臨界點時,就需要根據(jù)供需關(guān)系適當(dāng)作出調(diào)整。比如,當(dāng)廠家做了一大堆東西,產(chǎn)能過剩時,應(yīng)該暫停生產(chǎn),擴大宣傳,讓消費者過來消費。

同理,當(dāng)消費者發(fā)現(xiàn)某個熱銷商品售罄,應(yīng)該提醒廠家盡快生產(chǎn)。

在上面的案例中,生產(chǎn)者和消費者是不同種類的線程,一個負(fù)責(zé)存入,另一個負(fù)責(zé)取出,且它們操作的是同一個資源。但最難的部分在于:資源到達上限時,生產(chǎn)者等待,消費者消費;資源達到下限時,生產(chǎn)者生產(chǎn),消費者等待。
我們可以發(fā)現(xiàn),原本互不打擾的兩個線程之間開始了 “溝通”:
- 生產(chǎn)者:做的商品太多了,應(yīng)該擴大宣傳,讓大家來買。
- 消費者:都賣完啦,應(yīng)當(dāng)提醒商家盡快補貨。
這種線程間的相互調(diào)度,也就是線程間通信。
2、線程間通信的實現(xiàn)
實現(xiàn)線程間通信的方式有多種:
- 輪詢:生產(chǎn)者和消費者線程通過循環(huán)不斷檢查隊列的狀態(tài)。這種方法簡單,但會消耗大量 CPU 資源,且無法保證原子性。
- 等待喚醒機制(wait/notify):通過
wait和notify機制,線程可以在隊列為空或滿時阻塞自己,當(dāng)狀態(tài)改變時由其他線程喚醒。synchronized保證了線程的原子性,但notify可能導(dǎo)致線程競爭不均。 - 等待喚醒機制(Condition):使用
ReentrantLock和Condition實現(xiàn)等待喚醒機制,可以更加精確地控制線程的阻塞和喚醒。通過創(chuàng)建不同的Condition實例,可以分別管理生產(chǎn)者和消費者的等待狀態(tài),避免了notify的隨機喚醒問題。
2.1、輪詢
設(shè)計理念:生產(chǎn)者和消費者線程通過循環(huán)不斷檢查隊列的狀態(tài),隊列為空時生產(chǎn)者才可插入數(shù)據(jù),隊列不為空時消費者才能取出數(shù)據(jù),否則一律 sleep 等待。

代碼實現(xiàn):
import java.util.LinkedList;
import java.util.concurrent.TimeUnit;
/**
* 自定義阻塞隊列實現(xiàn):輪詢版本
*
* @param <T> 隊列中存儲的元素類型
*/
public class WhileQueue<T> {
// 用來存儲元素的容器
private final LinkedList<T> queue = new LinkedList<>();
// 隊列的最大容量
private final int MAX_SIZE = 1;
/**
* 將元素添加到隊列中
*
* @param resource 要插入的元素
* @throws InterruptedException 如果當(dāng)前線程被中斷
*/
public void put(T resource) throws InterruptedException {
// 如果隊列滿了,生產(chǎn)者線程將進入輪詢等待狀態(tài)
while (queue.size() >= MAX_SIZE) {
System.out.println("生產(chǎn)者:隊列已滿,無法插入...");
TimeUnit.MILLISECONDS.sleep(1000); // 線程等待1秒鐘再重試
}
// 插入元素到隊列的前面
System.out.println("生產(chǎn)者:插入" + resource + "!!!");
queue.addFirst(resource);
}
/**
* 從隊列中取出元素
*
* @throws InterruptedException 如果當(dāng)前線程被中斷
*/
public void take() throws InterruptedException {
// 如果隊列為空,消費者線程將進入輪詢等待狀態(tài)
while (queue.size() <= 0) {
System.out.println("消費者:隊列為空,無法取出...");
TimeUnit.MILLISECONDS.sleep(1000); // 線程等待1秒鐘再重試
}
// 從隊列的末尾取出元素
System.out.println("消費者:取出消息!!!");
queue.removeLast();
TimeUnit.MILLISECONDS.sleep(5000); // 模擬消費操作需要時間
}
}測試:
/**
* 測試類:創(chuàng)建生產(chǎn)者和消費者線程來測試WhileQueue的功能
*/
public class Test {
public static void main(String[] args) {
// 創(chuàng)建一個WhileQueue實例
WhileQueue<String> queue = new WhileQueue<>();
// 創(chuàng)建并啟動生產(chǎn)者線程
new Thread(new Runnable() {
@Override
public void run() {
for (int i = 0; i < 100; i++) {
try {
queue.put("消息" + i); // 插入消息到隊列
} catch (InterruptedException e) {
e.printStackTrace(); // 捕獲并打印中斷異常
}
}
}
}).start();
// 創(chuàng)建并啟動消費者線程
new Thread(new Runnable() {
@Override
public void run() {
for (int i = 0; i < 100; i++) {
try {
queue.take(); // 從隊列中取出消息
} catch (InterruptedException e) {
e.printStackTrace(); // 捕獲并打印中斷異常
}
}
}
}).start();
}
}由于設(shè)定了隊列最多只能存1個消息,所以只有當(dāng)隊列為空時,生產(chǎn)者才能插入數(shù)據(jù)。這是最簡單的線程間通信:多個線程不斷輪詢共享資源,通過共享資源的狀態(tài)判斷自己下一步該做什么。
但上面的實現(xiàn)方式存在一些缺點:
- 輪詢的方式太耗費 CPU 資源,如果線程過多,比如幾百上千個線程同時在那輪詢,會給 CPU 帶來較大負(fù)擔(dān)
- 無法保證原子性(代碼里沒有演示,但理論上確實如此,如果生產(chǎn)者的操作非原子性,消費者極可能獲取到臟數(shù)據(jù))
2.2、等待喚醒機制(wait/notify)
相對而言,等待喚醒機制則要優(yōu)雅得多,底層維護線程隊列,線程可以在隊列為空或滿時阻塞自己,當(dāng)狀態(tài)改變時由其他線程喚醒。synchronized 保證了線程的原子性,同時避免了過多線程同時自旋造成的 CPU 資源浪費,頗有點用空間換時間的味道。
當(dāng)一個生產(chǎn)者線程無法插入數(shù)據(jù)時,就讓它在隊列里休眠(阻塞),此時生產(chǎn)者線程會釋放 CPU 資源,等到消費者搶到 CPU 執(zhí)行權(quán)并取出數(shù)據(jù)后,再由消費者喚醒生產(chǎn)者繼續(xù)生產(chǎn)。
Java 有多種方式可以實現(xiàn)等待喚醒機制,最經(jīng)典的就是通過 wait 和 notify 的方式:
import java.util.LinkedList;
/**
* 自定義阻塞隊列實現(xiàn):使用 wait/notify
*
* @param <T> 隊列中存儲的元素類型
*/
public class WaitNotifyQueue<T> {
// 用來存儲元素的容器
private final LinkedList<T> queue = new LinkedList<>();
// 隊列的最大容量
private final int MAX_SIZE = 1;
/**
* 將元素添加到隊列中
*
* @param resource 要插入的元素
* @throws InterruptedException 如果當(dāng)前線程被中斷
*/
public synchronized void put(T resource) throws InterruptedException {
// 當(dāng)隊列滿時,生產(chǎn)者線程進入等待狀態(tài)
while (queue.size() >= MAX_SIZE) {
System.out.println("生產(chǎn)者:隊列已滿,無法插入...");
this.wait(); // 釋放鎖,并進入等待狀態(tài)
}
// 插入元素到隊列的前面
System.out.println("生產(chǎn)者:插入" + resource + "!!!");
queue.addFirst(resource);
this.notify(); // 喚醒等待的消費者線程
}
/**
* 從隊列中取出元素
*
* @throws InterruptedException 如果當(dāng)前線程被中斷
*/
public synchronized void take() throws InterruptedException {
// 當(dāng)隊列為空時,消費者線程進入等待狀態(tài)
while (queue.size() <= 0) {
System.out.println("消費者:隊列為空,無法取出...");
this.wait(); // 釋放鎖,并進入等待狀態(tài)
}
// 從隊列的末尾取出元素
System.out.println("消費者:取出消息!!!");
queue.removeLast();
this.notify(); // 喚醒等待的生產(chǎn)者線程
}
}基于 wait 和 notify 的阻塞隊列。其原理是通過同步機制和線程通信來處理生產(chǎn)者-消費者問題。在 put 方法中,生產(chǎn)者線程檢查隊列是否已滿,如果已滿,則調(diào)用 wait 使自己進入等待狀態(tài),釋放鎖,直到隊列有空位。生產(chǎn)者在插入元素后調(diào)用 notify 喚醒可能等待的消費者線程。在 take 方法中,消費者線程檢查隊列是否為空,如果為空,則調(diào)用 wait 使自己進入等待狀態(tài),釋放鎖,直到隊列有新元素。消費者在取出元素后調(diào)用 notify 喚醒可能等待的生產(chǎn)者線程。這種機制避免了忙等待,通過有效的線程通信提高了資源利用效率。
Ps:使用 notifyAll 在某些情況下可能更合適,尤其是當(dāng)有多個生產(chǎn)者和消費者線程時。notifyAll 會喚醒所有等待的線程,而不僅僅是一個線程,這樣可以保證系統(tǒng)中的所有線程都有機會被喚醒,避免了因線程喚醒不充分導(dǎo)致的潛在問題。
2.3、等待喚醒機制(Condition)
等待喚醒機制(wait/notify)版本的缺點是隨機喚醒容易出現(xiàn)"己方喚醒己方",最終導(dǎo)致全部線程阻塞的烏龍事件,雖然 wait/notifyAll 能解決這個問題,但喚醒全部線程又不夠精確,會造成無謂的線程競爭(實際只需要喚醒敵方線程即可)。
因此使用ReentrantLock和Condition實現(xiàn)等待喚醒機制,可以更加精確地控制線程的阻塞和喚醒。通過創(chuàng)建不同的Condition實例,可以分別管理生產(chǎn)者和消費者的等待狀態(tài),避免了notify的隨機喚醒問題。
作為改進版,可以使用 ReentrantLock 的 Condition 替代 synchronized 和 wait/notify:
import java.util.LinkedList;
import java.util.concurrent.locks.Condition;
import java.util.concurrent.locks.ReentrantLock;
public class ConditionQueue<T> {
// 容器,用來裝東西
private final LinkedList<T> queue = new LinkedList<>();
private final int CAPACITY = 10; // 隊列容量
// 顯式鎖(相對地,synchronized鎖被稱為隱式鎖)
private final ReentrantLock lock = new ReentrantLock();
private final Condition producerCondition = lock.newCondition();
private final Condition consumerCondition = lock.newCondition();
public void put(T resource) throws InterruptedException {
lock.lock();
try {
while (queue.size() >= CAPACITY) {
// 隊列滿了,不能再塞東西了,等待消費者取出數(shù)據(jù)
System.out.println("生產(chǎn)者:隊列已滿,無法插入...");
// 生產(chǎn)者阻塞
producerCondition.await();
}
System.out.println("生產(chǎn)者:插入" + resource + "!!!");
queue.addFirst(resource);
// 生產(chǎn)完畢,喚醒消費者
consumerCondition.signal();
} finally {
lock.unlock();
}
}
public void take() throws InterruptedException {
lock.lock();
try {
while (queue.size() <= 0) {
// 隊列空了,不能再取東西,等待生產(chǎn)者插入數(shù)據(jù)
System.out.println("消費者:隊列為空,無法取出...");
// 消費者阻塞
consumerCondition.await();
}
System.out.println("消費者:取出消息!!!");
queue.removeLast();
// 消費完畢,喚醒生產(chǎn)者
producerCondition.signal();
} finally {
lock.unlock();
}
}
}如何理解 Condition 呢?可以認(rèn)為 lock.newCondition() 創(chuàng)建了一個隊列,調(diào)用 producerCondition.await() 會把生產(chǎn)者線程放入生產(chǎn)者的等待隊列中,當(dāng)消費者調(diào)用producerCondition.signal() 時會喚醒從生產(chǎn)者的等待隊列中喚醒一個生產(chǎn)者線程出來工作。
也就是說,ReentrantLock 的 Condition 通過拆分線程等待隊列,讓線程的等待喚醒更加精確了,想喚醒哪一方就喚醒哪一方。
3、自定義阻塞隊列
基于以上機制,我們可以自定義實現(xiàn)一個簡單的阻塞隊列。以下代碼示例展示了一個基于 wait/notifyAll 實現(xiàn)的阻塞隊列:
public class BlockingQueue<T> {
private final LinkedList<T> queue = new LinkedList<>();
private int MAX_SIZE = 1;
private int remainCount = 0;
public BlockingQueue(int capacity) {
if (capacity <= 0) {
throw new IllegalArgumentException("size最小為1");
}
this.MAX_SIZE = capacity;
}
public synchronized void put(T resource) throws InterruptedException {
while (queue.size() >= MAX_SIZE) {
this.wait();
}
queue.addFirst(resource);
remainCount++;
this.notifyAll();
}
public synchronized T take() throws InterruptedException {
while (queue.size() <= 0) {
this.wait();
}
T resource = queue.removeLast();
remainCount--;
this.notifyAll();
return resource;
}
}4、Java 中的 BlockingQueue
BlockingQueue 是 Java 并發(fā)包(java.util.concurrent)中的一個接口,繼承自 Queue 接口。它提供了額外的阻塞操作,例如在隊列為空時等待元素變得可用,或在隊列已滿時等待空間變得可用。
BlockingQueue 阻塞隊列在 Java 中的主要實現(xiàn)有三個:
ArrayBlockingQueue: 基于數(shù)組實現(xiàn)的有界阻塞隊列,必須指定固定容量,支持可選的公平性策略。LinkedBlockingQueue: 基于鏈表實現(xiàn)的阻塞隊列,默認(rèn)無界或指定容量,有較高的插入和刪除性能。SynchronousQueue: 一個沒有內(nèi)部容量的隊列,每個插入操作必須等待一個對應(yīng)的刪除操作,反之亦然,適用于直接交換數(shù)據(jù)的場景。
更多實現(xiàn)可以參考:Java 并發(fā)集合:阻塞隊列集合介紹
到此這篇關(guān)于如何理解 Java 中的阻塞隊列:從基礎(chǔ)到高級的深度解析的文章就介紹到這了,更多相關(guān)java阻塞隊列內(nèi)容請搜索腳本之家以前的文章或繼續(xù)瀏覽下面的相關(guān)文章希望大家以后多多支持腳本之家!
相關(guān)文章
MyBatis基于pagehelper實現(xiàn)分頁原理及代碼實例
這篇文章主要介紹了MyBatis基于pagehelper實現(xiàn)分頁原理及代碼實例,文中通過示例代碼介紹的非常詳細(xì),對大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價值,需要的朋友可以參考下2020-06-06
Spring Boot 事務(wù)實戰(zhàn)之如何解決 DB 與 MQ
本文結(jié)合代碼案例探討如何利用Spring的 TransactionSynchronizationManager 實現(xiàn)事務(wù)提交后觸發(fā)(Trigger After Commit)機制,優(yōu)雅解決數(shù)據(jù)庫與消息隊列的"雙寫一致性"問題,感興趣的朋友跟隨小編一起看看吧2025-12-12
Java之注解@Data和@ToString(callSuper=true)解讀
在使用Lombok庫的@Data注解時,若子類未通過@ToString(callSuper=true)注明包含父類屬性,toString()方法只打印子類屬性,解決方法:1. 子類重寫toString方法;2. 子類使用@Data和@ToString(callSuper=true),父類也應(yīng)使用@Data2024-11-11
使用springboot跳轉(zhuǎn)到指定頁面和(重定向,請求轉(zhuǎn)發(fā)的實例)
這篇文章主要介紹了使用springboot跳轉(zhuǎn)到指定頁面和(重定向,請求轉(zhuǎn)發(fā)的實例),具有很好的參考價值,希望對大家有所幫助。如有錯誤或未考慮完全的地方,望不吝賜教2021-12-12
Java 網(wǎng)絡(luò)爬蟲基礎(chǔ)知識入門解析
這篇文章主要介紹了Java 網(wǎng)絡(luò)爬蟲基礎(chǔ)知識入門解析,文中通過示例代碼介紹的非常詳細(xì),對大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價值,需要的朋友可以參考下2019-10-10

