暂无图片
暂无图片
暂无图片
暂无图片
暂无图片

Java多线程进阶(三十)—— J.U.C之collections框架:ConcurrentLinkedDeque

TPVLOG 2021-06-21
257


本文首发于Ressmix个人站点:https://www.tpvlog.com

一、引言

在开始讲 ConcurrentLinkedDeque
之前,我们先来了解下Deque这种数据结构,我们知道Queue是一种具有FIFO特点的数据结构,元素只能在队首进行“入队”操作,在队尾进行“出队”操作。

Deque(double-ended queue)是一种双端队列,也就是说可以在任意一端进行“入队”,也可以在任意一端进行“出队”:

  1. public interface Deque<E> extends Queue<E>

Deque的数据结构示意图如下: 

我们再来看下JDK中QueueDeque这两种数据结构的接口定义,看看Deque和Queue相比有哪些增强:

1.1 Queue接口定义

Queue的接口非常简单,一共只有三种类型的操作:入队、出队、读取。

上述方法,可以划分如下:

操作类型抛出异常返回特殊值
入队add(e)offer(e)
出队remove()poll()
读取element()peek()

每种操作类型,都给出了两种方法,区别就是其中一种操作在队列的状态不满足某些要求时,会抛出异常;另一种,则直接返回特殊值(如null)。

1.2 Deque接口定义

Queue接口的所有方法Deque都具备,只不过队首/队尾都可以进行“出队”和“入队”操作:

操作类型抛出异常返回特殊值
队首入队addFirst(e)offerFirst(e)
队首出队removeFirst()pollFirst()
队首读取getFirst()peekFirst()
队尾入队addLast(e)offerLast(e)
队尾出队removeLast()pollLast()
队尾读取getLast()peekLast()

除此之外,Deque还可以当作“栈”来使用,我们知道“栈”是一种具有“LIFO”特点的数据结构(关于栈,可以参考我的这篇博文:栈),Deque提供了 push
、 pop
、 peek
这三个栈方法,一般实现这三个方法时,可以利用已有方法,即有如下映射关系:

栈方法Deque方法
pushaddFirst(e)
popremoveFirst()
peekpeekFirst()

关于Deque接口的更多细节,读者可以参考Oracle的官方文档:https://docs.oracle.com/javase/8/docs/api/

二、ConcurrentLinkedDeque简介

ConcurrentLinkedDeque
是JDK1.7时,J.U.C包引入的一个集合工具类。在JDK1.7之前,除了Stack类外,并没有其它适合并发环境的“栈”数据结构。ConcurrentLinkedDeque作为双端队列,可以当作“栈”来使用,并且高效地支持并发环境。

ConcurrentLinkedDequ
和 ConcurrentLinkedQueue
一样,采用了无锁算法,底层基于自旋+CAS的方式实现。

  1. public class ConcurrentLinkedDeque<E>

  2. extends AbstractCollection<E>

  3. implements Deque<E>, java.io.Serializable

三、ConcurrentLinkedDeque原理

3.1 队列结构

我们先来看下ConcurrentLinkedDeque的内部结构:

  1. public class ConcurrentLinkedDeque<E> extends AbstractCollection<E>

  2. implements Deque<E>, java.io.Serializable {


  3. /**

  4. * 头指针

  5. */

  6. private transient volatile Node<E> head;


  7. /**

  8. * 尾指针

  9. */

  10. private transient volatile Node<E> tail;


  11. private static final Node<Object> PREV_TERMINATOR, NEXT_TERMINATOR;


  12. // Unsafe mechanics

  13. private static final sun.misc.Unsafe UNSAFE;

  14. private static final long headOffset;

  15. private static final long tailOffset;


  16. static {

  17. PREV_TERMINATOR = new Node<Object>();

  18. PREV_TERMINATOR.next = PREV_TERMINATOR;

  19. NEXT_TERMINATOR = new Node<Object>();

  20. NEXT_TERMINATOR.prev = NEXT_TERMINATOR;

  21. try {

  22. UNSAFE = sun.misc.Unsafe.getUnsafe();

  23. Class<?> k = ConcurrentLinkedDeque.class;

  24. headOffset = UNSAFE.objectFieldOffset(k.getDeclaredField("head"));

  25. tailOffset = UNSAFE.objectFieldOffset(k.getDeclaredField("tail"));

  26. } catch (Exception e) {

  27. throw new Error(e);

  28. }

  29. }


  30. /**

  31. * 双链表结点定义

  32. */

  33. static final class Node<E> {

  34. volatile Node<E> prev; // 前驱指针

  35. volatile E item; // 结点值

  36. volatile Node<E> next; // 后驱指针


  37. Node() {

  38. }


  39. Node(E item) {

  40. UNSAFE.putObject(this, itemOffset, item);

  41. }


  42. boolean casItem(E cmp, E val) {

  43. return UNSAFE.compareAndSwapObject(this, itemOffset, cmp, val);

  44. }


  45. void lazySetNext(Node<E> val) {

  46. UNSAFE.putOrderedObject(this, nextOffset, val);

  47. }


  48. boolean casNext(Node<E> cmp, Node<E> val) {

  49. return UNSAFE.compareAndSwapObject(this, nextOffset, cmp, val);

  50. }


  51. void lazySetPrev(Node<E> val) {

  52. UNSAFE.putOrderedObject(this, prevOffset, val);

  53. }


  54. boolean casPrev(Node<E> cmp, Node<E> val) {

  55. return UNSAFE.compareAndSwapObject(this, prevOffset, cmp, val);

  56. }


  57. // Unsafe mechanics


  58. private static final sun.misc.Unsafe UNSAFE;

  59. private static final long prevOffset;

  60. private static final long itemOffset;

  61. private static final long nextOffset;


  62. static {

  63. try {

  64. UNSAFE = sun.misc.Unsafe.getUnsafe();

  65. Class<?> k = Node.class;

  66. prevOffset = UNSAFE.objectFieldOffset(k.getDeclaredField("prev"));

  67. itemOffset = UNSAFE.objectFieldOffset(k.getDeclaredField("item"));

  68. nextOffset = UNSAFE.objectFieldOffset(k.getDeclaredField("next"));

  69. } catch (Exception e) {

  70. throw new Error(e);

  71. }

  72. }

  73. }


  74. // ...

  75. }

可以看到,ConcurrentLinkedDeque的内部和ConcurrentLinkedQueue类似,不过是一个双链表结构,每入队一个元素就是插入一个Node类型的结点。字段 head
指向队列头, tail
指向队列尾,通过Unsafe来CAS操作字段值以及Node对象的字段值。

需要特别注意的是ConcurrentLinkedDeque包含两个特殊字段:PREVTERMINATOR、NEXTTERMINATOR。这两个字段初始时都指向一个值为null的空结点,这两个字段在结点删除时使用,后面会详细介绍:


3.2 构造器定义

ConcurrentLinkedDeque
包含两种构造器:

  1. /**

  2. * 空构造器.

  3. */

  4. public ConcurrentLinkedDeque() {

  5. head = tail = new Node<E>(null);


  6. }

  1. /**

  2. * 从已有集合,构造队列

  3. */

  4. public ConcurrentLinkedDeque(Collection<? extends E> c) {

  5. Node<E> h = null, t = null;

  6. for (E e : c) {

  7. checkNotNull(e);

  8. Node<E> newNode = new Node<E>(e);

  9. if (h == null)

  10. h = t = newNode;

  11. else { // 在队尾插入元素

  12. t.lazySetNext(newNode);

  13. newNode.lazySetPrev(t);

  14. t = newNode;

  15. }

  16. }

  17. initHeadTail(h, t);

  18. }

我们重点看下空构造器,通过空构造器建立的ConcurrentLinkedDeque对象,其head和tail指针并非指向null,而是指向一个item值为null的Node结点——哨兵结点,如下图:

3.3 入队操作

双端队列与普通队列的入队区别是:双端队列既可以在“队尾”插入元素,也可以在“队首”插入元素。ConcurrentLinkedDeque的入队方法有很多: addFirst(e)
、 addLast(e)
、 offerFirst(e)
、 offerLast(e)

  1. public void addFirst(E e) {

  2. linkFirst(e);

  3. }


  4. public void addLast(E e) {

  5. linkLast(e);

  6. }


  7. public boolean offerFirst(E e) {

  8. linkFirst(e);

  9. return true;

  10. }


  11. public boolean offerLast(E e) {

  12. linkLast(e);

  13. return true;

  14. }

可以看到,队首“入队”其实就是调用了 linkFirst(e)
方法,而队尾“入队”是调用了 linkLast(e)方法。我们先来看下队首“入队”——linkFirst(e):

  1. /**

  2. * 在队首插入一个元素.

  3. */

  4. private void linkFirst(E e) {

  5. checkNotNull(e);

  6. final Node<E> newNode = new Node<E>(e); // 创建待插入的结点


  7. restartFromHead:

  8. for (; ; )

  9. for (Node<E> h = head, p = h, q; ; ) {

  10. if ((q = p.prev) != null && (q = (p = q).prev) != null)

  11. // Check for head updates every other hop.

  12. // If p == q, we are sure to follow head instead.

  13. p = (h != (h = head)) ? h : q;

  14. else if (p.next == p) // PREV_TERMINATOR

  15. continue restartFromHead;

  16. else {

  17. // p is first node

  18. newNode.lazySetNext(p); // CAS piggyback

  19. if (p.casPrev(null, newNode)) {

  20. // Successful CAS is the linearization point

  21. // for e to become an element of this deque,

  22. // and for newNode to become "live".

  23. if (p != h) // hop two nodes at a time

  24. casHead(h, newNode); // Failure is OK.

  25. return;

  26. }

  27. // Lost CAS race to another thread; re-read prev

  28. }

  29. }

  30. }


为了便于理解,我们以示例来看:假设有两个线程ThreadA和ThreadB同时进行入队操作。

①ThreadA先单独入队一个元素9

此时,ThreadA会执行CASE3分支:

  1. else { // CASE3: p是队首结点

  2. newNode.lazySetNext(p); // “新结点”的next指向队首结点

  3. if (p.casPrev(null, newNode)) { // 队首结点的prev指针指向“新结点”

  4. if (p != h) // hop two nodes at a time

  5. casHead(h, newNode); // Failure is OK.

  6. return;

  7. }

  8. // 执行到此处说明CAS操作失败,有其它线程也在队首插入元素

  9. }

队列的结构如下:

 


②ThreadA入队一个元素2,同时ThreadB入队一个元素10

此时,依然执行CASE3分支,我们假设ThreadA操作成功,ThreadB操作失败:

  1. else { // CASE3: p是队首结点

  2. newNode.lazySetNext(p); // “新结点”的next指向队首结点

  3. if (p.casPrev(null, newNode)) { // 队首结点的prev指针指向“新结点”

  4. if (p != h) // hop two nodes at a time

  5. casHead(h, newNode); // Failure is OK.

  6. return;

  7. }

  8. // 执行到此处说明CAS操作失败,有其它线程也在队首插入元素

  9. }

ThreadA的CAS操作成功后,会进入以下判断:

  1. if (p != h) // hop two nodes at a time

  2. casHead(h, newNode); // Failure is OK.

上述判断的作用就是重置head头指针,可以看到,ConcurrentLinkedDeque其实是以每次跳2个结点的方式移动指针,这主要考虑到并发环境以这种hop跳的方式可以提升效率。

此时队列的结构如下: 

注意,此时ThreadB的 p.casPrev(null,newNode)
操作失败了,所以会进入下一次自旋,在下一次自旋中继续进入CASE3。如果ThreadA的casHead操作没有完成,ThreadB就进入了下一次自旋,则会进入分支1,重置指针p指向队首。最终队列结构如下:

在队尾插入元素和队首类似,不再赘述,读者可以自己阅读源码。

3.4 出队操作

ConcurrentLinkedDeque的出队一样分为队首、队尾两种情况: removeFirst()
、 pollFirst()
、 removeLast()
、 pollLast()

  1. public E removeFirst() {

  2. return screenNullResult(pollFirst());

  3. }


  4. public E removeLast() {

  5. return screenNullResult(pollLast());

  6. }


  7. public E pollFirst() {

  8. for (Node<E> p = first(); p != null; p = succ(p)) {

  9. E item = p.item;

  10. if (item != null && p.casItem(item, null)) {

  11. unlink(p);

  12. return item;

  13. }

  14. }

  15. return null;

  16. }


  17. public E pollLast() {

  18. for (Node<E> p = last(); p != null; p = pred(p)) {

  19. E item = p.item;

  20. if (item != null && p.casItem(item, null)) {

  21. unlink(p);

  22. return item;

  23. }

  24. }

  25. return null;

  26. }

可以看到,两个remove方法其实内部都调用了对应的poll方法,我们重点看下队尾的“出队”——pollLast方法:

  1. public E pollLast() {

  2. for (Node<E> p = last(); p != null; p = pred(p)) {

  3. E item = p.item;

  4. if (item != null && p.casItem(item, null)) {

  5. unlink(p);

  6. return item;

  7. }

  8. }

  9. return null;

  10. }

last方法用于寻找队尾结点,即满足 p.next==null&&p.prev!=p
的结点:

  1. Node<E> last() {

  2. restartFromTail:

  3. for (; ; )

  4. for (Node<E> t = tail, p = t, q; ; ) {

  5. if ((q = p.next) != null &&

  6. (q = (p = q).next) != null)

  7. // Check for tail updates every other hop.

  8. // If p == q, we are sure to follow tail instead.

  9. p = (t != (t = tail)) ? t : q;

  10. else if (p == t

  11. // It is possible that p is NEXT_TERMINATOR,

  12. // but if so, the CAS is guaranteed to fail.

  13. || casTail(t, p))

  14. return p;

  15. else

  16. continue restartFromTail;

  17. }

  18. }

pred方法用于寻找当前结点的前驱结点(如果前驱是自身,则返回队尾结点):

  1. final Node<E> pred(Node<E> p) {

  2. Node<E> q = p.prev;

  3. return (p == q) ? last() : q;

  4. }

unlink方法断开结点的链接:

  1. /**

  2. * Unlinks non-null node x.

  3. */

  4. void unlink(Node<E> x) {

  5. // assert x != null;

  6. // assert x.item == null;

  7. // assert x != PREV_TERMINATOR;

  8. // assert x != NEXT_TERMINATOR;


  9. final Node<E> prev = x.prev;

  10. final Node<E> next = x.next;

  11. if (prev == null) {

  12. unlinkFirst(x, next);

  13. } else if (next == null) {

  14. unlinkLast(x, prev);

  15. } else {

  16. Node<E> activePred, activeSucc;

  17. boolean isFirst, isLast;

  18. int hops = 1;


  19. // Find active predecessor

  20. for (Node<E> p = prev; ; ++hops) {

  21. if (p.item != null) {

  22. activePred = p;

  23. isFirst = false;

  24. break;

  25. }

  26. Node<E> q = p.prev;

  27. if (q == null) {

  28. if (p.next == p)

  29. return;

  30. activePred = p;

  31. isFirst = true;

  32. break;

  33. } else if (p == q)

  34. return;

  35. else

  36. p = q;

  37. }


  38. // Find active successor

  39. for (Node<E> p = next; ; ++hops) {

  40. if (p.item != null) {

  41. activeSucc = p;

  42. isLast = false;

  43. break;

  44. }

  45. Node<E> q = p.next;

  46. if (q == null) {

  47. if (p.prev == p)

  48. return;

  49. activeSucc = p;

  50. isLast = true;

  51. break;

  52. } else if (p == q)

  53. return;

  54. else

  55. p = q;

  56. }


  57. // TODO: better HOP heuristics

  58. if (hops < HOPS

  59. // always squeeze out interior deleted nodes

  60. && (isFirst | isLast))

  61. return;


  62. // Squeeze out deleted nodes between activePred and

  63. // activeSucc, including x.

  64. skipDeletedSuccessors(activePred);

  65. skipDeletedPredecessors(activeSucc);


  66. // Try to gc-unlink, if possible

  67. if ((isFirst | isLast) &&


  68. // Recheck expected state of predecessor and successor

  69. (activePred.next == activeSucc) &&

  70. (activeSucc.prev == activePred) &&

  71. (isFirst ? activePred.prev == null : activePred.item != null) &&

  72. (isLast ? activeSucc.next == null : activeSucc.item != null)) {


  73. updateHead(); // Ensure x is not reachable from head

  74. updateTail(); // Ensure x is not reachable from tail


  75. // Finally, actually gc-unlink

  76. x.lazySetPrev(isFirst ? prevTerminator() : x);

  77. x.lazySetNext(isLast ? nextTerminator() : x);

  78. }

  79. }

  80. }

ConcurrentLinkedDeque相比ConcurrentLinkedQueue,功能更丰富,但是由于底层结构是双链表,且完全采用CAS+自旋的无锁算法保证线程安全性,所以需要考虑各种并发情况,源码比ConcurrentLinkedQueue更加难懂,留待有精力作进一步分析。

四、总结

ConcurrentLinkedDeque使用了自旋+CAS的非阻塞算法来保证线程并发访问时的数据一致性。由于队列本身是一种双链表结构,所以虽然算法看起来很简单,但其实需要考虑各种并发的情况,实现复杂度较高,并且ConcurrentLinkedDeque不具备实时的数据一致性,实际运用中,如果需要一种线程安全的栈结构,可以使用ConcurrentLinkedDeque。

另外,关于ConcurrentLinkedDeque还有以下需要注意的几点:

  1. ConcurrentLinkedDeque的迭代器是弱一致性的,这在并发容器中是比较普遍的现象,主要是指在一个线程在遍历队列结点而另一个线程尝试对某个队列结点进行修改的话不会抛出ConcurrentModificationException,这也就造成在遍历某个尚未被修改的结点时,在next方法返回时可以看到该结点的修改,但在遍历后再对该结点修改时就看不到这种变化。

  2. size方法需要遍历链表,所以在并发情况下,其结果不一定是准确的,只能供参考。

文章转载自TPVLOG,如果涉嫌侵权,请发送邮件至:contact@modb.pro进行举报,并提供相关证据,一经查实,墨天轮将立刻删除相关内容。

评论