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

Java并發(fā)LinkedBlockingQueue源碼分析

 更新時(shí)間:2023年02月12日 15:48:24   作者:歷河川  
這篇文章主要為大家介紹了Java并發(fā)LinkedBlockingQueue源碼分析,有需要的朋友可以借鑒參考下,希望能夠有所幫助,祝大家多多進(jìn)步,早日升職加薪

簡介

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ì),pollpeek都是非阻塞的,但是區(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)文章

最新評(píng)論

江山市| 仙桃市| 江陵县| 德州市| 鄂尔多斯市| 夏河县| 塘沽区| 佛教| 华蓥市| 綦江县| 常宁市| 梨树县| 定远县| 平度市| 长治市| 宜良县| 长春市| 瑞丽市| 珠海市| 嘉善县| 安丘市| 成武县| 杨浦区| 拜泉县| 浮梁县| 江北区| 焉耆| 高阳县| 龙江县| 青州市| 乌海市| 乌兰察布市| 三原县| 乌拉特后旗| 化德县| 井研县| 汾阳市| 青河县| 万源市| 红安县| 剑阁县|