Java并發(fā)LinkedBlockingQueue源碼分析
簡介
LinkedBlockingQueue是一個(gè)阻塞的有界隊(duì)列,底層是通過一個(gè)個(gè)的Node節(jié)點(diǎn)形成的鏈表實(shí)現(xiàn)的,鏈表隊(duì)列中的頭節(jié)點(diǎn)是一個(gè)空的Node節(jié)點(diǎn),在多線程下操作時(shí)會(huì)使用ReentrantLock鎖來保證數(shù)據(jù)的安全性,并使用ReentrantLock下的Condition對(duì)象來阻塞以及喚醒線程。
常量
/**
* 鏈表中的節(jié)點(diǎn)類
*/
static class Node<E> {
//節(jié)點(diǎn)中的元素
E item;
//下一個(gè)節(jié)點(diǎn)
Node<E> next;
Node(E x) { item = x; }
}
/** 鏈表隊(duì)列的容量大小,如果沒有指定則使用Integer最大值 */
private final int capacity;
/** 記錄鏈表中的節(jié)點(diǎn)的數(shù)量的原子類 */
private final AtomicInteger count = new AtomicInteger();
/**鏈表的頭節(jié)點(diǎn)
*/
transient Node<E> head;
/**
* 鏈表的尾節(jié)點(diǎn)
*/
private transient Node<E> last;
/** 從鏈表隊(duì)列中獲取節(jié)點(diǎn)時(shí)防止多個(gè)線程同時(shí)操作所產(chǎn)生數(shù)據(jù)安全問題時(shí)所加的鎖 */
private final ReentrantLock takeLock = new ReentrantLock();
/** Wait queue for waiting takes */
private final Condition notEmpty = takeLock.newCondition();
/** 添加節(jié)點(diǎn)到鏈表隊(duì)列中防止多個(gè)線程同時(shí)操作所產(chǎn)生數(shù)據(jù)安全問題時(shí)所加的鎖 */
private final ReentrantLock putLock = new ReentrantLock();
/** Wait queue for waiting puts */
private final Condition notFull = putLock.newCondition();
Node:鏈表隊(duì)列中的節(jié)點(diǎn),用于存放元素。capacity:鏈表隊(duì)列中最多能存放的節(jié)點(diǎn)數(shù)量,如果在創(chuàng)建LinkedBlockingQueue的時(shí)候沒有指定.則默認(rèn)最多存放的節(jié)點(diǎn)的數(shù)量為Integer的最大值。head:鏈表隊(duì)列中的頭節(jié)點(diǎn),一般來說頭節(jié)點(diǎn)都是一個(gè)沒有元素的空節(jié)點(diǎn)。last:鏈表隊(duì)列中的尾節(jié)點(diǎn)。takeLock:在獲取鏈表隊(duì)列中的節(jié)點(diǎn)的時(shí)候所加的鎖。putLock:在添加鏈表隊(duì)列中的節(jié)點(diǎn)的時(shí)候所加的鎖。Condition:當(dāng)線程需要進(jìn)行等待或者喚醒的時(shí)候則會(huì)調(diào)用該對(duì)象下的方法。
構(gòu)造方法
/**
* 創(chuàng)建默認(rèn)容量大小的鏈表隊(duì)列
*/
public LinkedBlockingQueue() {
this(Integer.MAX_VALUE);
}
/**
* 創(chuàng)建指定容量大小的鏈表隊(duì)列
*/
public LinkedBlockingQueue(int capacity) {
if (capacity <= 0) throw new IllegalArgumentException();
this.capacity = capacity;
//創(chuàng)建一個(gè)空節(jié)點(diǎn),并將該節(jié)點(diǎn)設(shè)置為頭尾節(jié)點(diǎn)
last = head = new Node<E>(null);
}
/**
* 根據(jù)指定集合中的元素創(chuàng)建一個(gè)默認(rèn)容量大小的鏈表隊(duì)列
*/
public LinkedBlockingQueue(Collection<? extends E> c) {
//創(chuàng)建默認(rèn)容量大小的鏈表隊(duì)列
this(Integer.MAX_VALUE);
//獲取添加元素節(jié)點(diǎn)的鎖
final ReentrantLock putLock = this.putLock;
//加鎖
putLock.lock();
try {
//鏈表中節(jié)點(diǎn)的數(shù)量
int n = 0;
//遍歷集合中的元素
for (E e : c) {
if (e == null)
throw new NullPointerException();
if (n == capacity)
throw new IllegalStateException("Queue full");
//為元素創(chuàng)建一個(gè)節(jié)點(diǎn),并將節(jié)點(diǎn)添加到鏈表的尾部,并設(shè)置節(jié)點(diǎn)為尾節(jié)點(diǎn)
enqueue(new Node<E>(e));
//鏈表中節(jié)點(diǎn)的數(shù)量自增
++n;
}
//記錄鏈表中節(jié)點(diǎn)的數(shù)量
count.set(n);
} finally {
//釋放鎖
putLock.unlock();
}
}
第一個(gè)和第三個(gè)構(gòu)造方法中都會(huì)調(diào)用第二個(gè)構(gòu)造方法,而在第二個(gè)構(gòu)造方法中會(huì)設(shè)置鏈表隊(duì)列中容納節(jié)點(diǎn)的數(shù)量以及創(chuàng)建一個(gè)空的頭節(jié)點(diǎn)來填充,再看第三個(gè)構(gòu)造方法中的代碼,首先會(huì)獲取putLock鎖,代表當(dāng)前是一個(gè)需要添加節(jié)點(diǎn)的線程,再將指定集合中的元素封裝成一個(gè)Node節(jié)點(diǎn),并依次將封裝的節(jié)點(diǎn)追加到鏈表隊(duì)列中的尾部,并使用AtomicInteger來記錄鏈表隊(duì)列中節(jié)點(diǎn)的數(shù)量。
put
public void put(E e) throws InterruptedException {
if (e == null) throw new NullPointerException();
int c = -1;
//為指定元素創(chuàng)建節(jié)點(diǎn)
Node<E> node = new Node<E>(e);
//獲取添加元素節(jié)點(diǎn)的鎖
final ReentrantLock putLock = this.putLock;
//獲取記錄鏈表節(jié)點(diǎn)數(shù)量的原子類
final AtomicInteger count = this.count;
//加鎖,如果加鎖的線程被中斷了則拋出異常
putLock.lockInterruptibly();
try {
//校驗(yàn)鏈表中的節(jié)點(diǎn)數(shù)量是否到達(dá)了指定的容量
//如果到達(dá)了指定的容量就進(jìn)行阻塞等待
//如果線程被喚醒了,但是鏈表中的節(jié)點(diǎn)數(shù)量還是未改變,則繼續(xù)阻塞等待
//只有當(dāng)頭節(jié)點(diǎn)出隊(duì),新的節(jié)點(diǎn)才能繼續(xù)添加
while (count.get() == capacity) {
notFull.await();
}
//將新節(jié)點(diǎn)添加到鏈表的尾部并設(shè)置為尾節(jié)點(diǎn)
enqueue(node);
//獲取沒有添加當(dāng)前節(jié)點(diǎn)時(shí)鏈表中的節(jié)點(diǎn)數(shù)量
//并更新鏈表中的節(jié)點(diǎn)數(shù)量
c = count.getAndIncrement();
if (c + 1 < capacity)
//喚醒等待添加節(jié)點(diǎn)的線程
//可能當(dāng)前線程在等待隊(duì)列中等待的時(shí)候
//有新的線程要執(zhí)行添加節(jié)點(diǎn)的操作
//但是鏈表的容量已經(jīng)到達(dá)最大,所以新的線程也會(huì)進(jìn)行等待
//當(dāng)前線程被喚醒了并且鏈表的容量沒有到達(dá)最大則嘗試去喚醒等待的線程
notFull.signal();
} finally {
//釋放鎖
putLock.unlock();
}
if (c == 0)
//c等于0說明添加當(dāng)前節(jié)點(diǎn)的時(shí)候鏈表中沒有節(jié)點(diǎn)
//可能有線程在獲取節(jié)點(diǎn),但是鏈表中沒有節(jié)點(diǎn)
//從而一直進(jìn)行等待,當(dāng)添加了節(jié)點(diǎn)的時(shí)候就需要喚醒獲取節(jié)點(diǎn)的線程
signalNotEmpty();
}
LinkedBlockingQueue中的代碼都比較簡單,主要是ReentrantLock下的Condition中的方法比較復(fù)雜,我們先整體的了解一下put方法,首先通過new Node為將指定元素封裝成一個(gè)節(jié)點(diǎn),再獲取putLock鎖,當(dāng)鏈表隊(duì)列中的節(jié)點(diǎn)數(shù)量已經(jīng)到達(dá)了capacity大小,那當(dāng)前線程就需要調(diào)用Condition下的await方法進(jìn)行等待將線程阻塞,直到有節(jié)點(diǎn)出隊(duì)或者說有節(jié)點(diǎn)被刪除或者當(dāng)前線程被中斷了,當(dāng)前線程被中斷了則會(huì)直接退出當(dāng)前put方法并拋出異常,如果節(jié)點(diǎn)出隊(duì)了或者節(jié)點(diǎn)被刪除了,那當(dāng)前線程被喚醒了則會(huì)繼續(xù)執(zhí)行添加節(jié)點(diǎn)的操作。
enqueue方法則會(huì)將封裝的節(jié)點(diǎn)追加到鏈表隊(duì)列中的尾部,通過getAndIncrement方法先獲取沒有添加當(dāng)前節(jié)點(diǎn)時(shí)鏈表隊(duì)列中節(jié)點(diǎn)的數(shù)量,然后更新添加了當(dāng)前節(jié)點(diǎn)之后鏈表隊(duì)列中節(jié)點(diǎn)的數(shù)量,c則是沒有添加當(dāng)前節(jié)點(diǎn)時(shí)鏈表隊(duì)列中節(jié)點(diǎn)的數(shù)量,c+1則是添加當(dāng)前節(jié)點(diǎn)后鏈表隊(duì)列中節(jié)點(diǎn)的數(shù)量,如果說c+1小于capacity則說明線程在添加節(jié)點(diǎn)的時(shí)候,鏈表隊(duì)列中的節(jié)點(diǎn)數(shù)量已經(jīng)到達(dá)了最大值,后續(xù)添加節(jié)點(diǎn)的線程都需要進(jìn)行阻塞,當(dāng)有節(jié)點(diǎn)被刪除或出隊(duì)的時(shí)候,最開始阻塞的線程被喚醒,被喚醒的線程則會(huì)去執(zhí)行添加節(jié)點(diǎn)的操作,當(dāng)添加完節(jié)點(diǎn)之后鏈表隊(duì)列中的節(jié)點(diǎn)數(shù)量沒有到達(dá)最大值則會(huì)去喚醒后續(xù)被阻塞的線程執(zhí)行添加節(jié)點(diǎn)的操作。
c等于0說明在添加當(dāng)前節(jié)點(diǎn)之前,可能有線程在獲取鏈表隊(duì)列中的節(jié)點(diǎn),但是鏈表隊(duì)列中沒有節(jié)點(diǎn),導(dǎo)致獲取節(jié)點(diǎn)的線程處于阻塞狀態(tài),當(dāng)添加完節(jié)點(diǎn)之后,鏈表隊(duì)列中有了節(jié)點(diǎn),此時(shí)就需要喚醒阻塞的線程去獲取節(jié)點(diǎn)。
添加元素的方法分為put和offer,區(qū)別在于阻塞與非阻塞,當(dāng)鏈表隊(duì)列中的節(jié)點(diǎn)數(shù)量已經(jīng)到達(dá)最大值,put方法則會(huì)阻塞,而offer方法不會(huì)阻塞則是直接返回。
獲取元素的方法分為take、poll、peek,take方法與put方法相似,只不過一個(gè)是入隊(duì),一個(gè)是出隊(duì),poll與peek都是非阻塞的,但是區(qū)別在于poll獲取了節(jié)點(diǎn)之后,該節(jié)點(diǎn)會(huì)從鏈表隊(duì)列中移除,而peek不會(huì)移除節(jié)點(diǎn)。
await
public final void await() throws InterruptedException {
if (Thread.interrupted())
//線程被中斷拋出異常
throw new InterruptedException();
//為當(dāng)前線程創(chuàng)建一個(gè)等待模式的節(jié)點(diǎn)并入隊(duì),并將等待隊(duì)列中已經(jīng)取消等待的節(jié)點(diǎn)移除掉
Node node = addConditionWaiter();
//釋放當(dāng)前線程的鎖,防止當(dāng)前線程加了鎖,導(dǎo)致其它在等待的線程被喚醒之后不能獲取到鎖從而導(dǎo)致一直阻塞
int savedState = fullyRelease(node);
int interruptMode = 0;
//如果指定節(jié)點(diǎn)還在等待隊(duì)列中等待則掛起
//如果指定節(jié)點(diǎn)被中斷了則會(huì)將指定節(jié)點(diǎn)添加到同步等待隊(duì)列中
//如果指定節(jié)點(diǎn)被喚醒了則會(huì)將指定節(jié)點(diǎn)添加到同步等待隊(duì)列中
while (!isOnSyncQueue(node)) {
//節(jié)點(diǎn)在等待隊(duì)列中則掛起
LockSupport.park(this);
//線程在等待隊(duì)列中被中斷則會(huì)添加到同步等待隊(duì)列中
if ((interruptMode = checkInterruptWhileWaiting(node)) != 0)
break;
}
//acquireQueued 指定節(jié)點(diǎn)中的線程被中斷了或者被喚醒了則會(huì)嘗試去獲取鎖
//如果還未到指定節(jié)點(diǎn)中的線程獲取鎖的時(shí)候則會(huì)繼續(xù)掛起
if (acquireQueued(node, savedState) && interruptMode != THROW_IE)
interruptMode = REINTERRUPT;
if (node.nextWaiter != null)
//指定節(jié)點(diǎn)的線程已經(jīng)獲取到了鎖并且節(jié)點(diǎn)關(guān)聯(lián)的下一個(gè)節(jié)點(diǎn)不為空
//此時(shí)就需要將已經(jīng)獲取到鎖的節(jié)點(diǎn)從等待隊(duì)列中移除
unlinkCancelledWaiters();
if (interruptMode != 0)
reportInterruptAfterWait(interruptMode);
}
首先通過addConditionWaiter方法將當(dāng)前線程封裝成一個(gè)等待模式的節(jié)點(diǎn),并將節(jié)點(diǎn)添加到等待隊(duì)列中以及會(huì)將等待隊(duì)列中已經(jīng)取消等待的線程節(jié)點(diǎn)從隊(duì)列中移除,再通過fullyRelease方法釋放掉當(dāng)前線程加的所有的鎖,之所以釋放鎖是防止其它線程獲取不到鎖從而一直阻塞,再看isOnSyncQueue方法,該方法是校驗(yàn)當(dāng)前線程節(jié)點(diǎn)是否在等待隊(duì)列中,如果在等待隊(duì)列中那就將節(jié)點(diǎn)中的線程掛起等待。
isOnSyncQueue
final boolean isOnSyncQueue(Node node) {
if (node.waitStatus == Node.CONDITION || node.prev == null)
//指定節(jié)點(diǎn)還在等待隊(duì)列中此時(shí)就需要繼續(xù)等待
return false;
if (node.next != null)
//指定節(jié)點(diǎn)已經(jīng)不在等待隊(duì)列中了
return true;
//從等待隊(duì)列中的尾節(jié)點(diǎn)開始向頭節(jié)點(diǎn)遍歷,校驗(yàn)指定的節(jié)點(diǎn)是否在其中
return findNodeFromTail(node);
}
當(dāng)節(jié)點(diǎn)的狀態(tài)為CONDITION時(shí),則說明該節(jié)點(diǎn)還在等待隊(duì)列中,node.prev等于null為什么說也是在等待隊(duì)列中呢?因?yàn)榈却?duì)列中的節(jié)點(diǎn)是沒有prev指針和next指針的,如果prev指針和next指針指向的節(jié)點(diǎn)不為空,那就說明該節(jié)點(diǎn)是在同步等待隊(duì)列中的,如果在同步等待隊(duì)列中的話,那節(jié)點(diǎn)中的線程就可以嘗試去獲取鎖并執(zhí)行后續(xù)的操作。
當(dāng)?shù)却?duì)列中的線程節(jié)點(diǎn)被喚醒和中斷則會(huì)添加到同步等待隊(duì)列中,如果是被中斷的話則會(huì)通過checkInterruptWhileWaiting方法添加一個(gè)中斷標(biāo)識(shí),再通過acquireQueued方法來獲取鎖,如果獲取鎖失敗則繼續(xù)等待,當(dāng)獲取鎖成功之后則會(huì)該節(jié)點(diǎn)從等待隊(duì)列中移除,如果說你是一個(gè)被中斷的線程,最后會(huì)通過reportInterruptAfterWait方法拋出中斷異常。
signal
public final void signal() {
if (!isHeldExclusively())
//加鎖的線程不是當(dāng)前線程則拋出異常
throw new IllegalMonitorStateException();
//頭節(jié)點(diǎn)
Node first = firstWaiter;
if (first != null)
//喚醒頭節(jié)點(diǎn)
doSignal(first);
}
/**
* 喚醒等待隊(duì)列中的頭節(jié)點(diǎn)
* 如果等待隊(duì)列中的頭節(jié)點(diǎn)被取消等待或已經(jīng)被喚醒了
* 此時(shí)就需要喚醒頭節(jié)點(diǎn)的后續(xù)的一個(gè)節(jié)點(diǎn)
* 直到成功的喚醒一個(gè)節(jié)點(diǎn)中的線程
*/
private void doSignal(Node first) {
do {
if ( (firstWaiter = first.nextWaiter) == null)
lastWaiter = null;
first.nextWaiter = null;
} while (!transferForSignal(first) && (first = firstWaiter) != null);
}
/**
* 將指定的節(jié)點(diǎn)添加到同步等待隊(duì)列中
* 并根據(jù)前一個(gè)節(jié)點(diǎn)的等待狀態(tài)來決定是否需要立刻喚醒指定節(jié)點(diǎn)
*/
final boolean transferForSignal(Node node) {
if (!compareAndSetWaitStatus(node, Node.CONDITION, 0))
//更改節(jié)點(diǎn)狀態(tài)失敗說明該節(jié)點(diǎn)已經(jīng)被喚醒了
return false;
//將要喚醒的節(jié)點(diǎn)添加到同步等待隊(duì)列中
//并返回前一個(gè)節(jié)點(diǎn)
Node p = enq(node);
//前一個(gè)節(jié)點(diǎn)的等待狀態(tài)
int ws = p.waitStatus;
if (ws > 0 || !compareAndSetWaitStatus(p, ws, Node.SIGNAL))
//如果前一個(gè)節(jié)點(diǎn)的等待狀態(tài)大于0則說明已經(jīng)被取消加鎖,此時(shí)就需要喚醒后續(xù)的節(jié)點(diǎn),就是當(dāng)前節(jié)點(diǎn)
//前一個(gè)節(jié)點(diǎn)的等待狀態(tài)不大于0但是更改前一個(gè)節(jié)點(diǎn)的等待狀態(tài)時(shí)失敗則說明前一個(gè)節(jié)點(diǎn)已經(jīng)被喚醒了并更改了狀態(tài)
//此時(shí)就需要嘗試將當(dāng)前節(jié)點(diǎn)中的線程喚醒
LockSupport.unpark(node.thread);
return true;
}
喚醒線程節(jié)點(diǎn)的方法主要還是看transferForSignal方法,首先會(huì)通過cas操作將需要喚醒的節(jié)點(diǎn)的狀態(tài)設(shè)置為0,如果更改節(jié)點(diǎn)狀態(tài)失敗則說明該節(jié)點(diǎn)已經(jīng)被喚醒了,更新節(jié)點(diǎn)狀態(tài)成功則會(huì)通過enq方法將節(jié)點(diǎn)添加到同步等待隊(duì)列中,此時(shí)就需要根據(jù)前一個(gè)節(jié)點(diǎn)來決定是否需要立即喚醒當(dāng)前節(jié)點(diǎn)中的線程。
從下面的圖片中能看出來其實(shí)同步等待隊(duì)列和等待隊(duì)列中使用的節(jié)點(diǎn)是共用的節(jié)點(diǎn),并不會(huì)創(chuàng)建新的節(jié)點(diǎn),同步等待隊(duì)列中的節(jié)點(diǎn)使用next指針和prev指針來關(guān)聯(lián)節(jié)點(diǎn),而等待隊(duì)列中則是使用nextWaiter指針來關(guān)聯(lián)節(jié)點(diǎn)的。

以上就是Java并發(fā)LinkedBlockingQueue源碼分析的詳細(xì)內(nèi)容,更多關(guān)于Java并發(fā)LinkedBlockingQueue的資料請(qǐng)關(guān)注腳本之家其它相關(guān)文章!
相關(guān)文章
關(guān)于Controller 層返回值的公共包裝類的問題
本文給大家介紹Controller 層返回值的公共包裝類-避免每次都包裝一次返回-InitializingBean增強(qiáng),本文通過實(shí)例代碼給大家介紹的非常詳細(xì),需要的朋友參考下吧2021-09-09
Java數(shù)據(jù)結(jié)構(gòu)及算法實(shí)例:三角數(shù)字
這篇文章主要介紹了Java數(shù)據(jù)結(jié)構(gòu)及算法實(shí)例:三角數(shù)字,本文直接給出實(shí)現(xiàn)代碼,代碼中包含詳細(xì)注釋,需要的朋友可以參考下2015-06-06
springboot接口多實(shí)現(xiàn)類選擇性注入解決方案
這篇文章主要為大家介紹了springboot接口多實(shí)現(xiàn)類選擇性注入解決方案的四種方式,有需要的朋友可以借鑒參考下,希望能夠有所幫助,祝大家多多進(jìn)步2022-03-03
Spring MVC前后端的數(shù)據(jù)傳輸?shù)膶?shí)現(xiàn)方法
這篇文章主要介紹了Spring MVC前后端的數(shù)據(jù)傳輸?shù)膶?shí)現(xiàn)方法,需要的朋友可以參考下2017-10-10
Eclipse中如何引入JUnit進(jìn)行單元測(cè)試
這篇文章主要介紹了Eclipse中如何引入JUnit進(jìn)行單元測(cè)試問題,具有很好的參考價(jià)值,希望對(duì)大家有所幫助,如有錯(cuò)誤或未考慮完全的地方,望不吝賜教2024-04-04
Springboot 2.x集成kafka 2.2.0的示例代碼
kafka近幾年更新非??欤部梢钥闯鰇afka在企業(yè)中是用的頻率越來越高。本文主要為大家介紹了Springboot 2.x集成kafka 2.2.0的示例代碼,需要的可以參考一下2022-04-04

