锁和并发队列
锁和并发队列
一、ConcurrentLinkedQueue
ConcurrentLinkedQueue是Java提供的线程安全无边界非阻塞队列,其底层数据结构使用链表实现。对于入队和出队操作使用CAS非阻塞原子操作来实现线程安全。
1、代码结构(jdk21)
ConcurrentLinkedQueue内部使用单向链表来实现队列功能,其中内部有两个Node类型的变量用来指向链表的头节点和尾节点,这两个变量使用volatile关键字修饰保证了其内存可见性,通过无参构造方法我们可以看到默认情况下一个初始化的队列,链表的头尾节点都指向一个指为null的哨兵节点。当执行入队操作时元素被添加到队列尾部,出队时从队列头取队列元素。
ConcurrentLinkedQueue.java
public class ConcurrentLinkedQueue<E> extends AbstractQueue<E> implements Queue<E>, java.io.Serializable { private static final long serialVersionUID = 196745693267521676L; transient volatile Node<E> head; private transient volatile Node<E> tail; /** * Creates a {@code ConcurrentLinkedQueue} that is initially empty. */ public ConcurrentLinkedQueue() { head = tail = new Node<E>(); } // ......}下面我们来看一下表示节点的Node类,这是一个ConcurrentLinkedQueue.java的内部类,内部有两个用volatile修饰的属性item和next,item表示节点的值,next表示下一个节点也就是队列的下一个元素。从这里可以看出来只要有足够的内存,链表就可以一直往下添加,所以这是一个无界的单向列表。从代码(3)可以看到,对于volatile变量item、next的操作使用的是VarHandle类提供的CAS算法来保证线程安全。VorHandle.java是jdk9提供的用来替代Unsafe类实现CAS算法的工具类。
ConcurrentLinkedQueue.java
public class ConcurrentLinkedQueue<E> extends AbstractQueue<E> implements Queue<E>, java.io.Serializable { // ...... static final class Node<E> { volatile E item; // (1) volatile Node<E> next; // (2) /** * Constructs a node holding item. Uses relaxed write because * item can only be seen after piggy-backing publication via CAS. */ Node(E item) { ITEM.set(this, item); } /** Constructs a dead dummy node. */ Node() {} void appendRelaxed(Node<E> next) { // assert next != null; // assert this.next == null; NEXT.set(this, next); } boolean casItem(E cmp, E val) { // assert item == cmp || item == null; // assert cmp != null; // assert val == null; return ITEM.compareAndSet(this, cmp, val); } } // VarHandle mechanics (3) private static final VarHandle HEAD; private static final VarHandle TAIL; static final VarHandle ITEM; static final VarHandle NEXT; static { try { MethodHandles.Lookup l = MethodHandles.lookup(); HEAD = l.findVarHandle(ConcurrentLinkedQueue.class, "head", Node.class); TAIL = l.findVarHandle(ConcurrentLinkedQueue.class, "tail", Node.class); ITEM = l.findVarHandle(Node.class, "item", Object.class); NEXT = l.findVarHandle(Node.class, "next", Node.class); } catch (ReflectiveOperationException e) { throw new ExceptionInInitializerError(e); } } // ......}2、offer方法原理
boolean offer(E e)方法向队列添加一个元素,如果添加成功返回true,如果参数e是null,则抛出NIP异常。从源码看这个方法并不会返回false,除非e为null抛出异常或者执行成功返回true,否则方法会一直循环竞争CAS资源。因为使用的是CAS无阻塞算法,这个方法不会阻塞线程。
public boolean offer(E e) { final Node<E> newNode = new Node<E>(Objects.requireNonNull(e)); // (1) for (Node<E> t = tail, p = t;;) { Node<E> q = p.next; // (2) if (q == null) { // p is last node if (NEXT.compareAndSet(p, null, newNode)) { // (3) if (p != t) // hop two nodes at a time; failure is OK TAIL.weakCompareAndSet(this, t, newNode); // (4) return true; } // Lost CAS race to another thread; re-read next } else if (p == q) p = (t != (t = tail)) ? t : head; // (5) else // Check for tail updates after two hops. p = (p != t && t != (t = tail)) ? t : q; // (6) }}我们重点看代码(3)的实现,这是offer方法实现原子操作保证线程安全的关键,NEXT.compareAndSet(p, null, newNode)表示如果p对象的的next为null则把p的next设置newNode返回true,否则说明在执行代码之前已经有其他线程修改了p.next的值,方法返回false,线程开始下一次循环直到NEXT.compareAndSet(p, null, newNode)执行成功返回true。这是由CAS算法实现的非阻塞原子性操作,具体由VarHandle工具类实现。
3、poll方法原理
E poll()方法取队列头部的元素,然后将该元素从队列里删除,如果队列为空则返回null。
public E poll() { restartFromHead: for (;;) { // (1) for (Node<E> h = head, p = h, q;; p = q) { // (2) final E item; if ((item = p.item) != null && p.casItem(item, null)) { // (3) if (p != h) // hop two nodes at a time updateHead(h, ((q = p.next) != null) ? q : p); // (4) return item; } else if ((q = p.next) == null) { updateHead(h, p); // (5) return null; } else if (p == q) continue restartFromHead; // (6) } }}代码(1),定义了一个goto标记,从队列头开始执行循环
代码(2),找到链表的头节点
代码(3),取出头节点的值item,如果值不为null,则返回取到的值,在返回之前需要使用CAS操作把当前节点的item值设置为null标记为删除状态,表示该节点已被删除。
代码(4),把更新头节点为新的节点,原来的头节点将会被删除(没有对象引用后会被gc回收)。
代码(6),回到goto标记restartFromHead从链表头开始循环。
和offer方法一样,poll方法也是使用CAS非阻塞原子算法保证更新队列操作的线程安全,p.casItem(item, null)表示把对象p的item值等于参数item,则把对象p的item更新为null并返回true,否则返回false。
boolean casItem(E cmp, E val) { // assert item == cmp || item == null; // assert cmp != null; // assert val == null; return ITEM.compareAndSet(this, cmp, val);}static final VarHandle ITEM;MethodHandles.Lookup l = MethodHandles.lookup();ITEM = l.findVarHandle(Node.class, "item", Object.class);ITEM是一个VarHandle对象,通过它可以对Node对象的item进行CAS比较交换操作。通过非阻塞原子操作实现线程安全。
4、peek方法原理
E peek()方法返回队列的第一个元素,如果队列为空则返回null,和poll方法不同的是peek方法不会删除元素。
public E peek() { restartFromHead: for (;;) { for (Node<E> h = head, p = h, q;; p = q) { final E item; if ((item = p.item) != null || (q = p.next) == null) { updateHead(h, p); return item; } else if (p == q) continue restartFromHead; } }}peek方法和poll方法类似,从队列头开始循环,找到队列头节点并返回节点的值item,然后把队列的头节点设置为下一个节点。
5、size方法原理
int size()方法返回队列的当前元素个数,由于ConcurrentLinkedQueue队列使用CAS实现的线程安全,CAS没有使用锁,所以在size方法执行过程中可能有其他线程添加或删除了队列元素,所以在多线程环境下size方法的返回结果并不准确。
public int size() { restartFromHead: for (;;) { int count = 0; for (Node<E> p = first(); p != null;) { if (p.item != null) if (++count == Integer.MAX_VALUE) break; // @see Collection.size() if (p == (p = p.next)) continue restartFromHead; } return count; }}
加载评论中...