NOTE
6.41 PriorityBlockingQueue
What PriorityBlockingQueue is, binary-heap basics, usage, and source analysis of growth, offer/put, take, and heap adjustment.
This is a historical learning note and may contain outdated or incomplete understanding.
1. What Is It?
An unbounded blocking queue implemented underneath with an array representing a binary heap.
Reads and writes share one main lock, so queue operations are mutually exclusive while holding that lock.
Elements can be ordered.
Because the queue is unbounded from the API perspective, put does not block for capacity, while take blocks when the queue is empty.
1.1. Binary Heap
A binary heap is a complete binary tree whose heap-order property here is that each node’s value is less than or equal to the values of its left and right children. Therefore, the minimum value is at the root.
It is stored in an array. For element a[i], the left child is a[2*i+1], the right child is a[2*i+2], and the parent is a[(i-1)/2].
The structure is shown below:

2. How to Use It
public class PriorityBlockingQueueTest
{
public static void main(String[] args) throws InterruptedException
{
PriorityBlockingQueue<String> queue = new PriorityBlockingQueue<>(1);
CountDownLatch latch = new CountDownLatch(2);
new Thread(()->{
for (int i = 0;;i++)
{
try
{
String data = "data" + i;
queue.put(data);
System.out.println("Producer puts message: " + data);
TimeUnit.SECONDS.sleep(1);
}
catch (Exception e)
{
e.printStackTrace();
}
finally
{
latch.countDown();
}
}
}).start();
new Thread(()->{
for (;;)
{
try
{
System.out.println("Consumer gets message: " + queue.take());
}
catch (Exception e)
{
e.printStackTrace();
}
finally
{
latch.countDown();
}
}
}).start();
latch.await();
}
}
3. Implementation Analysis
3.1. Constructor
3.1.1. Implemented with Array + Lock + Condition
public class PriorityBlockingQueue<E> extends AbstractQueue<E>
implements BlockingQueue<E>, java.io.Serializable {
// Implemented underneath with an array (heap).
private transient Object[] queue;
// Actual number of elements.
private transient int size;
// comparator determines element ordering; null means natural ordering.
private transient Comparator<? super E> comparator;
// One main lock means queue operations are mutually exclusive while held.
private final ReentrantLock lock;
// One condition because only the read side needs to wait for data.
// Wake readers when the queue becomes non-empty.
private final Condition notEmpty;
// Used as an additional CAS-based spin lock when growing the array.
private transient volatile int allocationSpinLock;
public PriorityBlockingQueue() {
// Default capacity 11, natural ordering.
this(DEFAULT_INITIAL_CAPACITY, null);
}
public PriorityBlockingQueue(int initialCapacity) {
this(initialCapacity, 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];
}
}
3.2. put
public void put(E e) {
// Delegate to offer.
offer(e); // never need to block
}
3.2.1. Delegate to offer; Capacity Does Not Block
offer
public boolean offer(E e) {
if (e == null)
throw new NullPointerException();
final ReentrantLock lock = this.lock;
// Lock.
lock.lock();
int n, cap;
Object[] array;
// If the number of elements >= array length, grow the array.
while ((n = size) >= (cap = (array = queue).length))
tryGrow(array, cap);
try {
Comparator<? super E> cmp = comparator;
// Natural order: put e at position n and sift upward.
if (cmp == null)
siftUpComparable(n, e, array);
else
siftUpUsingComparator(n, e, array, cmp);
size = n + 1;
// Wake a reader.
notEmpty.signal();
} finally {
lock.unlock();
}
return true;
}
- Line 6: acquire the lock. While the writer holds it, other queue operations using this lock cannot proceed.
- Lines 10-11: grow the backing array when needed.
- Lines 13-18: place the element at the end and restore heap order with a sift-up operation.
- Line 19: increment the logical queue size.
- Line 21: a reader may be blocked because the queue was empty, so signal it.
Because the logical queue is unbounded, writing does not wait for a notFull condition.
3.2.1.1. Lock
final ReentrantLock lock = this.lock;
lock.lock();
try {
//...
} finally {
lock.unlock();
}
3.2.1.2. Check Whether Expansion Is Needed
while ((n = size) >= (cap = (array = queue).length))
tryGrow(array, cap);
3.2.1.2.1. Grow When Necessary
tryGrow
private void tryGrow(Object[] array, int oldCap) {
// Release the main lock so readers are not blocked for the entire allocation phase.
lock.unlock(); // must release and then re-acquire main lock
Object[] newArray = null;
// allocationSpinLock == 0 means no other thread currently owns the growth spin lock.
if (allocationSpinLock == 0 &&
UNSAFE.compareAndSwapInt(this, allocationSpinLockOffset,
0, 1)) {
try {
// If old capacity < 64, grow faster: new = old + old + 2.
// Otherwise grow by about 1.5x.
int newCap = oldCap + ((oldCap < 64) ?
(oldCap + 2) : // grow faster if small
(oldCap >> 1));
// Overflow check.
if (newCap - MAX_ARRAY_SIZE > 0) { // possible overflow
int minCap = oldCap + 1;
if (minCap < 0 || minCap > MAX_ARRAY_SIZE)
throw new OutOfMemoryError();
newCap = MAX_ARRAY_SIZE;
}
// Allocate only if expansion is still needed and queue still points to array.
if (newCap > oldCap && queue == array)
newArray = new Object[newCap];
} finally {
// Release growth spin lock.
allocationSpinLock = 0;
}
}
// Another thread may be growing; yield the CPU.
if (newArray == null) // back off if another thread is allocating
Thread.yield();
// Re-acquire the main lock before replacing the backing array.
lock.lock();
// Copy the old array into the new one.
if (newArray != null && queue == array) {
queue = newArray;
System.arraycopy(array, 0, newArray, 0, oldCap);
}
}
3.2.1.3. Add the Element at the End of the Heap
if (cmp == null)
siftUpComparable(n, e, array);
3.2.1.3.1. Sift Up to Restore Heap Order
siftUpComparable
// Insert x at position k in heap array.
private static <T> void siftUpComparable(int k, T x, Object[] array) {
Comparable<? super T> key = (Comparable<? super T>) x;
// At most move up to root index 0.
while (k > 0) {
// Parent position: (k - 1) / 2.
int parent = (k - 1) >>> 1;
Object e = array[parent];
// If x is >= parent, heap order is satisfied.
if (key.compareTo((T) e) >= 0)
break;
// Otherwise move the parent down.
array[k] = e;
// Continue upward from the parent position.
k = parent;
}
array[k] = key;
}
3.2.1.3.2. Adjustment Diagram

3.3. take
public E take() throws InterruptedException {
final ReentrantLock lock = this.lock;
// Lock.
lock.lockInterruptibly();
E result;
try {
// If dequeue returns null, the queue is empty: block and wait.
while ( (result = dequeue()) == null)
notEmpty.await();
} finally {
// Unlock.
lock.unlock();
}
return result;
}
- Line 4: acquire the lock.
- Lines 8-9: dequeue. If the queue is empty, wait until a writer signals that the queue is non-empty.
3.3.1. Lock
final ReentrantLock lock = this.lock;
lock.lockInterruptibly();
try {
//...
} finally {
lock.unlock();
}
3.3.2. Block Until Dequeue Succeeds
while ( (result = dequeue()) == null)
notEmpty.await();
3.3.2.1. Dequeue Operation
dequeue
private E dequeue() {
// Return null if the queue is empty.
int n = size - 1;
if (n < 0)
return null;
else {
Object[] array = queue;
// Root at index 0 is the dequeued element.
E result = (E) array[0];
E x = (E) array[n];//Last array element.
array[n] = null;
Comparator<? super E> cmp = comparator;
if (cmp == null)
// Move x to root and restore heap order.
siftDownComparable(0, x, array, n);
else
siftDownUsingComparator(0, x, array, n, cmp);
size = n;
return result;
}
}
- Lines 7-12: remove the heap root and take the last element.
- Lines 14-18: restore heap order by sifting the last element downward.
3.3.2.1.1. Remove the Heap Root and Move the Last Element to the Top
Object[] array = queue;
E result = (E) array[0];
E x = (E) array[n];
array[n] = null;
3.3.2.1.2. Sift Down to Restore Heap Order
siftDownComparable
private static <T> void siftDownComparable(int k, T x, Object[] array,
int n) {
if (n > 0) {
Comparable<? super T> key = (Comparable<? super T>)x;
// Only non-leaf nodes have children to compare.
int half = n >>> 1;
while (k < half) {
// Left child.
int child = (k << 1) + 1;
Object c = array[child];
// Right child.
int right = child + 1;
if (right < n &&
((Comparable<? super T>) c).compareTo((T) array[right]) > 0)
c = array[child = right];
// c is the smaller child.
if (key.compareTo((T) c) <= 0)
break;
// Move the smaller child upward.
array[k] = c;
// Continue downward.
k = child;
}
array[k] = key;
}
}
3.3.2.1.3. Adjustment Diagram

4. Summary
An unbounded priority queue implemented with a binary heap and therefore ordered by natural order or a comparator.
Writes do not block for capacity; reads block when the queue is empty.
Discussion
Sign in with GitHub to comment. Discussions are stored as GitHub Issues.View on GitHub