文章导读:
【原创】JUC包源码分析01 | ArrayBlockingQueue
【原创】JUC包源码分析02 | LinkedBlockingQueue
备注:JDK版本:1.8.0,系统环境:Ubuntu16.04 LTS
这两个月比较忙,文章没有更新,在此给自己敲敲警钟,学习无止境!
之前的文章中分别对JUC包下的ArrayBlockingQueue和LinkedBlockingQueue进行了源码分析,本文将对另一个常用的阻塞队列java.util.concurrent.PriorityBlockingQueue做详细的分析。
java.util.concurrent.PriorityBlockingQueue是一个无界的阻塞队列。无界,就意味着可以一直往这个队列里面插入数据,因此,该队列在使用的过程中可能会导致OOM的情况发生。使用该队列有一个硬性的规定:插入该队列的元素必须是可比较(Comparable)的元素,否则将不被允许插入,并抛出java.lang.ClassCastException异常。
java.util.concurrent.PriorityBlockingQueue中插入的元素的优先级相同的时候,可以自定义java.util.Comparators进行排序的确定。
java.util.concurrent.PriorityBlockingQueue类似于ArrayBlockingQueue、LinkedBlockingQueue,继承于AbstractQueue,实现了BlockingQueue接口(关于AbstractQueue和BlockingQueue的详细讲解,欢迎阅读我的另一篇文章-【原创】JUC包源码分析01 | ArrayBlockingQueue)。
public class PriorityBlockingQueue<E> extends AbstractQueue<E>implements BlockingQueue<E>, java.io.Serializable
本文将直接源码剖析java.util.concurrent.PriorityBlockingQueue。
public class PriorityBlockingQueue<E> extends AbstractQueue<E>implements BlockingQueue<E>, java.io.Serializable {private static final long serialVersionUID = 5595510919245408276L;/*** 默认的数组容量.*/private static final int DEFAULT_INITIAL_CAPACITY = 11;/*** 最大的数组容量.* OutOfMemoryError: 当数组容量超过内存限制后会抛出该异常*/private static final int MAX_ARRAY_SIZE = Integer.MAX_VALUE - 8;/*** Priorityueue表现为一个平衡的二进制堆: queue[n]的子级孩是* queue[2*n+1] 和 queue[2*(n+1)].* PriorityQueue借助comparator进行排序, 或者是元素的自然排序规则* 进行排序 (如果没有指定排序的comparator)* 最小的element是queue[0],并且这个queue不能为空*/private transient Object[] queue;/*** PriorityQueue的元素数量.*/private transient int size;/*** 比较器, 为null时则采用元素的自然排序.*/private transient Comparator<? super E> comparator;/*** Lock锁*/private final ReentrantLock lock;/*** 等待队列. 因为PriorityBlockingQueue是无界的,故不会存在满的情况。* 但会存在元素为空的情况,所以只需要一个notEmpty的等待队列*/private final Condition notEmpty;/*** 通过CAS获取的自旋锁.*/private transient volatile int allocationSpinLock;
/*** 创建默认容量(11)大小的PriorityBlockingQueue.* 采用元素的自然排序规则进行排序.*/public PriorityBlockingQueue() {this(DEFAULT_INITIAL_CAPACITY, null);}/*** 借助自定义的initialCapacity创建PriorityBlockingQueue* 采用元素的自然排序规则进行排序.*/public PriorityBlockingQueue(int initialCapacity) {this(initialCapacity, null);}/*** 借助自定义的initialCapacity和比较器comparator创建.* 如果comparator为null, 则采用元素的原始排序方式*/public PriorityBlockingQueue(int initialCapacity,Comparator<? super E> comparator) {if (initialCapacity < 1)throw new IllegalArgumentException();this.lock = new ReentrantLock();this.notEmpty = lock.newCondition();this.comparator = comparator;this.queue = new Object[initialCapacity];}/*** 借助collection对象创建PriorityBlockingQueue* 如果这个collection对象是SortedSet或者PriorityQueue,* 将借助这个collection对象的排序方式进行排序.* 否则,采用元素的自然排序.** @param c collection对象* @throws ClassCastException 如果元素不能被比较,将抛出这个异常* @throws NullPointerException 元素为null,将抛出这个异常*/public PriorityBlockingQueue(Collection<? extends E> c) {this.lock = new ReentrantLock();this.notEmpty = lock.newCondition();boolean heapify = true; // true if not known to be in heap orderboolean screen = true; // true if must screen for nullsif (c instanceof SortedSet<?>) {SortedSet<? extends E> ss = (SortedSet<? extends E>) c;this.comparator = (Comparator<? super E>) ss.comparator();heapify = false;}else if (c instanceof PriorityBlockingQueue<?>) {PriorityBlockingQueue<? extends E> pq =(PriorityBlockingQueue<? extends E>) c;this.comparator = (Comparator<? super E>) pq.comparator();screen = false;if (pq.getClass() == PriorityBlockingQueue.class) // exact matchheapify = false;}Object[] a = c.toArray();int n = a.length;// If c.toArray incorrectly doesn't return Object[], copy it.if (a.getClass() != Object[].class)a = Arrays.copyOf(a, n, Object[].class);if (screen && (n == 1 || this.comparator != null)) {for (int i = 0; i < n; ++i)if (a[i] == null)throw new NullPointerException();}this.queue = a;this.size = n;if (heapify)heapify();}
3、PriorityBlockingQueue插入方法源码分析
PriorityBlockingQueue的元素插入方法分为add()、offer()和put()方法,下面看这些方法的源码。
/*** 扩充queue的容量* 必须获取锁对象** @param array queue的数组对象* @param oldCap queue的数组对象old容量(扩充前的容量)*/private void tryGrow(Object[] array, int oldCap) {lock.unlock(); // must release and then re-acquire main lockObject[] newArray = null;if (allocationSpinLock == 0 &&// 借助UNSAFE类,用于实现CAS操作,// 保证只能有一个线程进行queue容量的扩充UNSAFE.compareAndSwapInt(this, allocationSpinLockOffset,0, 1)) {try {// 如果queue的元素数组的大小小于64,直接扩充2个,变成2n+2// 如果queue的容量大于64,则扩充为n+0.5n=1.5nint newCap = oldCap + ((oldCap < 64) ?(oldCap + 2) : // grow faster if small(oldCap >> 1));// 判断是否超过容量限制if (newCap - MAX_ARRAY_SIZE > 0) { // possible overflowint minCap = oldCap + 1;if (minCap < 0 || minCap > MAX_ARRAY_SIZE)throw new OutOfMemoryError(); // 超出抛出异常newCap = MAX_ARRAY_SIZE;}if (newCap > oldCap && queue == array)// 初始化新的array数组newArray = new Object[newCap];} finally {// 因为此时只是一个单独的线程在操作,因为不需要cas操作allocationSpinLock = 0;}}if (newArray == null) // back off if another thread is allocatingThread.yield(); //lock.lock(); // 获取锁if (newArray != null && queue == array) {queue = newArray;// 复制数据,与Arrays.copyOf()方法相比,更快System.arraycopy(array, 0, newArray, 0, oldCap);}}/*** 在位置k插入元素x* @param k 元素放置的位置* @param x 元素对象* @param array queue内部数组*/private static <T> void siftUpComparable(int k, T x, Object[] array) {// 不使用自定义的Comparator,直接显示的借用元素自身的排序方式Comparable<? super T> key = (Comparable<? super T>) x;while (k > 0) {// 采用类似二分法的方式// 循环比较待插入元素与原有元素的大小int parent = (k - 1) >>> 1;Object e = array[parent];if (key.compareTo((T) e) >= 0)break;array[k] = e;k = parent;}// 将k位置的元素插入到array中array[k] = key;}/*** 基于自定义的Comparator进行排序元素后在位置k插入元素*/private static <T> void siftUpUsingComparator(int k, T x, Object[] array,Comparator<? super T> cmp) {while (k > 0) {int parent = (k - 1) >>> 1;Object e = array[parent];// 借用自定义的Comparator进行比较if (cmp.compare(x, (T) e) >= 0)break;array[k] = e;k = parent;}array[k] = x;}/*** 插入元素.* @param e 待插入的元素* @return 插入成功返回true* @throws ClassCastException 当元素不能被比较时(comparable)抛出* @throws NullPointerException 元素为null,抛出*/public boolean add(E e) {return offer(e);}/*** 插入元素到队列中.* 因为该队列是无界的,因此,该方法将永远不返回false.* @param e 插入队列* @return 插入成功,返回true* @throws ClassCastException 如果元素不能与原有插入的元素进行比较,抛出* @throws NullPointerException 元素为null,抛出*/public boolean offer(E e) {if (e == null) // 元素为null,抛出异常throw new NullPointerException();final ReentrantLock lock = this.lock;lock.lock(); // 加锁int n, cap;Object[] array;while ((n = size) >= (cap = (array = queue).length))tryGrow(array, cap); // 队列元素数量不少于数组容量时,队列扩容try {Comparator<? super E> cmp = comparator;if (cmp == null)// 如果 Comparator 为null,则调用元素原生的排序规则siftUpComparable(n, e, array);else// 使用自定义的Comparator,进行排序siftUpUsingComparator(n, e, array, cmp);size = n + 1;// 既然已经有元素插入,队列自然不为空,因此可以唤醒等待的notEmpty等待队列notEmpty.signal();} finally {lock.unlock(); // 解锁}return true; // 因为该队列是无界的,因此,插入后返回true}/*** 插入元素到队列中.* 因为该队列的无界性,因此,该方法将不被阻塞.* @param e 待插入的元素* @return 插入成功返回true* @throws ClassCastException 当元素不能被比较时(comparable)抛出* @throws NullPointerException 元素为null,抛出*/public void put(E e) {offer(e); // never need to block}
从上面的代码注释中,可以发现PriorityBlockingQueue的插入元素方法与ArrayBlockingQueue和LinkedBlockingQueue有明显的区别,因为PriorityBlockingQueue是无界的,因此,不存在插入元素时队列满而发生阻塞插入的情况,因此正常情况下,只要是插入元素,都将插入成功(可能出现OOM异常)。add()方法和put()方法均调用的是无阻塞的offer()方法,用于将元素插入到队列中。
在元素插入的时候,会对元素进行排序,排序的规则借助于元素原生的排序方式或者是用户自定义的排序方式,以用户自定义的排序方式为准,当没有设定相应的排序方式时,则采用元素自身的排序方式。
4、PriorityBlockingQueue获取(删除)元素方法源码分析
PriorityBlockingQueue获取元素方法分为poll()、take()和peek()方法,删除元素的方法是remove()方法,poll()、take()方法和peek()方法相比,peek()方法只是获取元素,不删除队列中的元素,而poll()、take()、remove()方法将从队列中删除元素。
我们先看poll()、take()和peek()方法的源码。
/*** 从队列中获取元素* 必须在加锁状态*/private E dequeue() {int n = size - 1;if (n < 0) // 如果队列为空return null;else {Object[] array = queue;E result = (E) array[0]; // 获取数组的第一个元素E x = (E) array[n];array[n] = null;Comparator<? super E> cmp = comparator;if (cmp == null)// 移动元素,对数组进行排序siftDownComparable(0, x, array, n);else// 移动元素,对数组进行排序 (借助自定义的Comparator)siftDownUsingComparator(0, x, array, n, cmp);size = n;return result;}}/*** 队列中有元素,返回元素,队列为空,返回null*/public E poll() {final ReentrantLock lock = this.lock;lock.lock(); // 加锁try {// 获取锁对象return dequeue();} finally {lock.unlock(); // 释放锁}}/*** 从队列中获取元素,队列为空,将阻塞*/public E take() throws InterruptedException {final ReentrantLock lock = this.lock;lock.lockInterruptibly(); // 加锁,可以干扰的锁E result;try {// 循环获取队列中的元素,如果队列为空,则会一直阻塞等待while ( (result = dequeue()) == null)notEmpty.await();} finally {lock.unlock(); // 释放锁}return result;}/*** 获取队列的首元素,但是不删除元素*/public E peek() {final ReentrantLock lock = this.lock;lock.lock();try {return (size == 0) ? null : (E) queue[0];} finally {lock.unlock();}}
从上面的代码中,可以看出因为poll()、take()和peek()方法都是在加锁的状态下进行获取(删除)元素,因此,这几个方法都是线程安全的方法,不同点是poll()方法和peek()方法,在队列中没有元素的时候,直接返回null,take()方法会阻塞,当队列中添加了元素后,take线程会被唤醒,进而从队列中获取元素。
从代码注释中,我们可以发现poll()方法的执行逻辑是:在加锁的情况下,先判断队列是否有元素,没有就直接返回null,有元素的话,则获取数组的首元素,然后对剩下的元素进行排序。获取元素调用的是PriorityBlockingQueue#dequeue()方法。
take()方法的执行逻辑:加锁状态下,循环判断队列中是否有元素,队列中没有元素时,通过notEmpty条件队列的await()方法阻塞take线程,将线程放入等待队列中进行等待,当队列有元素时,插入线程唤醒等待的take线程,从而从队列中获取(移除)元素。
接下来,我们来分析PriorityBlockingQueue的remove()方法的源码。
/*** 如果队列中含有一个或者多个这个元素,将从队列中移除元素.* 判断的标准是o.equals(e)* 如果队列中含有这个元素,移除成功,返回true** @param o 待删除的元素* @return true 如果队列改变(移除成功),返回true*/public boolean remove(Object o) {final ReentrantLock lock = this.lock;lock.lock(); // 加锁try {int i = indexOf(o);if (i == -1)return false;removeAt(i); // 根据索引删除元素return true;} finally {lock.unlock(); // 解锁}}/*** 获取待删除元素的索引位置,判断的标准是o.equals(e)*/private int indexOf(Object o) {if (o != null) {Object[] array = queue;int n = size;for (int i = 0; i < n; i++)if (o.equals(array[i]))return i;}return -1;}
remove()方法,实际比较简单,在加锁的情况下,获取元素的索引位置,如果待删除的元素不存在,则返回false,否则,在该索引处删除对应的元素,并重新对数组进行排序操作。
本文选择性分析了PriorityBlockingQueue常用的核心方法,并不涵盖PriorityBlockingQueue的所有方法,感兴趣的朋友希望能够自行阅读源码并加以总结,后续将对JUC并发包里其他常用类做分析。





