概述

LinkedTransferQueue是单向链表结构的无界阻塞队列。它通过CAS和LockSupport来实现线程安全的。
它采用一种预占模式,意思是当消费者线程取元素时,如果队列为空,那就生成一个节点(节点元素为null)入队,然后消费者线程被等待在这个节点上,后面生产者线程入队时发现一个元素为null的节点,就会直接将元素填充到该节点,并唤醒该节点等待的线程,被唤醒的消费者线程取走元素并返回。

LinkedTransferQueue在实现上有几个特点:

  • 双重队列:LinkedTransferQueueNode存储了一个isData的字段,用来区分该节点代表的是数据还是请求,称为双重队列机制。
  • 松弛度:为了节省CAS操作的开销,LinkedTransferQueue为了节省CAS的开销,LinkedTransferQueue不会立即取更新head/tail,而是需要等到ehad/tail与最近一个未匹配节点之间的距离超过一个松弛度阈值的时候,才会更新。
  • 节点自链接:已匹配的节点的next引用会指向自身。如果在遍历时遇到一个自链节点,那就表明当前线程已经滞后于另外一个更新head的兴奋啊从,此时就需要重新获取head来遍历。

1pyD7d.png

源码分析

继承体系

1pyRc8.png

重要属性

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24

/** 队列头,第一次入列之前为空*/
transient volatile Node head;

/** 队列尾节点,第一次添加节点之前为空 */
private transient volatile Node tail;

/** 累计到一定的次数再清除无效的node */
private transient volatile int sweepVotes;

/** sweepVote的阈值 */
static final int SWEEP_THRESHOLD = 32;

/**当前驱节点正在处理,当前节点再阻塞之前的自旋次数*/
private static final int FRONT_SPINS = 1 << 7;

/**当前驱节点正在处理,当前节点再阻塞之前的自旋次数*/
private static final int CHAINED_SPINS = FRONT_SPINS >>> 1;


private static final int NOW = 0; // for untimed poll, tryTransfer
private static final int ASYNC = 1; // for offer, put, add
private static final int SYNC = 2; // for transfer, take
private static final int TIMED = 3; // for timed poll, tryTransfer

重要方法解析

xfer方法

入队/出队方法都是由xfer来实现,所以我们这里只对xfer进行解析

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
private E xfer(E e, boolean haveData, int how, long nanos) {
if (haveData && (e == null))
throw new NullPointerException();
Node s = null; // the node to append, if needed

retry:
for (;;) { // restart on append race

for (Node h = head, p = h; p != null;) { // find & match first node
//从head开始一直向后匹配
boolean isData = p.isData;
Object item = p.item;
if (item != p && (item != null) == isData) {//有效节点
if (isData == haveData)
//节点与此次操作模式一致,无法匹配
break;
if (p.casItem(item, e)) { //匹配成功,cas修改为指定元素
for (Node q = p; q != h;) {
Node n = q.next; // update by 2 unless singleton
if (head == h && casHead(h, n == null ? q : n)) { //更新head为匹配节点的next节点

//旧head节点指向自身等待回收
h.forgetNext();
break;
} // advance and retry
//cas失败,重新获取head
if ((h = head) == null ||
(q = h.next) == null || !q.isMatched())
//如果head的next节点未被匹配,跳出循环,不更新head,即松弛度小于2
break; // unless slack < 2
}
//唤醒节点上等待的线程
LockSupport.unpark(p.waiter);
return LinkedTransferQueue.<E>cast(item);
}
}
//匹配失败,继续向后查找节点
Node n = p.next;
p = (p != n) ? n : (h = head); // Use head if p offlist
}

//未找到匹配节点,吧当前节点假如到队列尾
if (how != NOW) { // No matches available
if (s == null)
s = new Node(e, haveData);
//将新节点s添加到队列尾并返回s的前驱节点
Node pred = tryAppend(s, haveData);
if (pred == null)
//与其它不同模式线程竞争是失败重新循环
continue retry; // lost race vs opposite mode
if (how != ASYNC)//同步操作,等待匹配
return awaitMatch(s, pred, e, (how == TIMED), nanos);
}
return e; // not waiting
}
}

xfer方法的基本执行流程:

  1. 从head开始往后遍历匹配,找到一个节点模式和本次操作不同的未匹配的节点进行匹配。
  2. 匹配成功的节点CAS修改匹配节点的item未给定的元素e
  3. 如果此时所匹配节点向后移动,则cas更新head节点为匹配节点的next节点,旧head节点连接指向自身等待被回收。如果cas失败,并且松弛度大于等于2,旧需要重新获取head。
  4. 匹配成功,唤醒匹配节点p的灯箱线程waiter,返回匹配的item。
  5. 如果在上述操作中没有找到匹配节点,则根据参数how的不同做不同的处理
  • NOW:立即返回
  • SYNC:通过tryAppend方法插入一个新的节点s到队列尾,然后自旋或阻塞当前线程知道节点被匹配或取消返回。
  • ASYNC:通过tryAppend方法插入一个新的节点s到队列尾,异步直接返回
  • TIMED:通过tryAppend反复插入一个新系欸但到队列尾,然后通过自旋或阻塞当前线程直到节点被匹配或取消或等待超时时返回。

tryAppend方法

尝试添加节点s作为尾节点

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
private Node tryAppend(Node s, boolean haveData) {
for (Node t = tail, p = t;;) { // move p to last node and append
Node n, u; // temps for reads of next & tail
if (p == null && (p = head) == null) { //链表未初始化
if (casHead(null, s)) //将s作为head节点
return s; // initialize
}
else if (p.cannotPrecede(haveData))
return null; // lost race vs opposite mode
else if ((n = p.next) != null) // not last; keep traversing
p = p != t && t != (u = tail) ? (t = u) : // stale tail
(p != n) ? n : null; // restart if off list
else if (!p.casNext(null, s))
p = p.next; // re-read on CAS failure
else {
if (p != t) { // update if slack now >= 2
while ((tail != t || !casTail(t, s)) &&
(t = tail) != null &&
(s = t.next) != null && // advance and retry
(s = s.next) != null && s != t);
}
return p;
}
}
}

该方法添加给定节点s到队列尾并返回s的前继节点,失败时(与其它不同模式线程竞争失败)返回null,没有前继节点返回自身。

awaitMatch方法

自旋/让步/阻塞,直到给定节点s匹配到或a放弃匹配。

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
private E awaitMatch(Node s, Node pred, E e, boolean timed, long nanos) {
final long deadline = timed ? System.nanoTime() + nanos : 0L;
Thread w = Thread.currentThread();
int spins = -1; // initialized after first item and cancel checks
ThreadLocalRandom randomYields = null; // bound if needed

for (;;) {
Object item = s.item;
if (item != e) { // matched
// assert item != s;
s.forgetContents(); // avoid garbage
return LinkedTransferQueue.<E>cast(item);
}
if ((w.isInterrupted() || (timed && nanos <= 0)) &&
s.casItem(e, s)) { //取消匹配,item指向自身
unsplice(pred, s); //解除s节点和前继节点的连接
return e;
}

if (spins < 0) { // establish spins at/near front
if ((spins = spinsFor(pred, s.isData)) > 0)
randomYields = ThreadLocalRandom.current();
}
else if (spins > 0) { // spin
--spins;
if (randomYields.nextInt(CHAINED_SPINS) == 0)
Thread.yield(); // occasionally yield
}
else if (s.waiter == null) {
s.waiter = w; // request unpark then recheck
}
else if (timed) {
nanos = deadline - System.nanoTime();
if (nanos > 0L)
LockSupport.parkNanos(this, nanos);
}
else {
LockSupport.park(this);
}
}
}

当前操作作为同步操作时,会调用awaitMatch方法阻塞等待匹配,成功返回匹配节点item,返回失败返回给定参数e。在等待期间如果线程被中断或等待超时,则取消p,并调用upsplice方法解除节点,和其前继节点的链接。

unsplice方法

解除给定已经被删除/取消节点和前继节点的链接,可能会延迟解除。

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
final void unsplice(Node pred, Node s) {
s.forgetContents(); // forget unneeded fields
/*
* See above for rationale. Briefly: if pred still points to
* s, try to unlink s. If s cannot be unlinked, because it is
* trailing node or pred might be unlinked, and neither pred
* nor s are head or offlist, add to sweepVotes, and if enough
* votes have accumulated, sweep.
*/
if (pred != null && pred != s && pred.next == s) {
Node n = s.next;
if (n == null ||
(n != s && pred.casNext(s, n) && pred.isMatched())) {
//解除s节点的链接
for (;;) { // check if at, or could be, head
Node h = head;
if (h == pred || h == s || h == null)
return; // at head or list empty
if (!h.isMatched())
break;
Node hn = h.next;
if (hn == null)
return; // now empty
if (hn != h && casHead(h, hn))//更新head
h.forgetNext(); // advance head
}
if (pred.next != pred && s.next != s) { // recheck if offlist
for (;;) { // sweep now if enough votes
int v = sweepVotes;
if (v < SWEEP_THRESHOLD) {
if (casSweepVotes(v, v + 1))
break;
}
else if (casSweepVotes(v, 0)) {
sweep();
break;
}
}
}
}
}
}

sweep

解除从头不遍历时遇到的已经被匹配的节点的链接

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15

private void sweep() {
for (Node p = head, s, n; p != null && (s = p.next) != null; ) {
if (!s.isMatched())
// Unmatched nodes are never self-linked
p = s;
else if ((n = s.next) == null) // trailing node is pinned
break;
else if (s == n) // stale
// No need to also check for p == s, since that implies s == n
p = head;
else
p.casNext(s, n);
}
}