NOTE
6.12 ArrayBlockingQueue
What ArrayBlockingQueue is, how to use it, and source analysis of its array, lock, conditions, put/take/offer/poll/add/remove/element/peek methods.
This is a historical learning note and may contain outdated or incomplete understanding.
1. What Is It?
A bounded blocking queue implemented with an Object array.
Reads and writes share one lock in this implementation, so competing operations are mutually exclusive while holding that lock.
2. How to Use It
public class ArrayBlockingQueueTest
{
public static void main(String[] args) throws InterruptedException
{
ArrayBlockingQueue<String> queue = new ArrayBlockingQueue<>(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();
}
}
2.1. Choosing Methods
| Method / Behavior | Throws Exception | Returns Special Value | Blocks | Times Out |
|---|---|---|---|---|
| Insert | add(e) | offer(e) | put(e) | offer(e,time,unit) |
| Remove | remove() | poll() | take() | poll(time,unit) |
| Inspect | element() | peek() | Not available | Not available |
3. Implementation Analysis
3.1. UML

3.2. Constructor
3.2.1. Implemented with Array + Lock + Condition
public class ArrayBlockingQueue<E> extends AbstractQueue<E>
implements BlockingQueue<E>, java.io.Serializable {
// Implemented underneath with an array.
final Object[] items;
// Position of the next element read by take, poll, peek, or remove.
int takeIndex;
// Position where put, offer, or add writes the next element.
int putIndex;
// Actual number of elements in the array.
// count == items.length means the array is full.
int count;
// One lock means reads and writes are mutually exclusive while holding it.
final ReentrantLock lock;
// Two conditions.
private final Condition notEmpty;//Used to wake readers.
private final Condition notFull;//Used to wake writers.
public ArrayBlockingQueue(int capacity, boolean fair) {
if (capacity <= 0)
throw new IllegalArgumentException();
this.items = new Object[capacity];
lock = new ReentrantLock(fair);
notEmpty = lock.newCondition();
notFull = lock.newCondition();
}
}
3.3. put [Blocking]
public void put(E e) throws InterruptedException {
checkNotNull(e);
// Lock.
final ReentrantLock lock = this.lock;
lock.lockInterruptibly();
try {
// If the array is full, wait. A reader wakes the writer after taking an element.
while (count == items.length)
notFull.await();
// Not full: enqueue.
enqueue(e);
} finally {
lock.unlock();
}
}
- Line 4: acquire the lock. Once the writer holds the lock, other readers/writers cannot enter the critical section at the same time.
- Lines 8-9: if the array is full, block and wait.
- Line 11: if not full, enqueue and wake readers.
Details follow.
3.3.1. Lock
final ReentrantLock lock = this.lock;
lock.lockInterruptibly();
try {
//...
} finally {
lock.unlock();
}
3.3.2. If the Array Is Full, Wait
while (count == items.length)
notFull.await();
3.3.3. If It Is Not Full, Enqueue and Wake Readers
enqueue(e);
- enqueue
private void enqueue(E x) {
// Add the element at the tail.
final Object[] items = this.items;
items[putIndex] = x;
// After inserting at the end, reset putIndex to 0.
// This array is reused circularly and does not need expansion.
if (++putIndex == items.length)
putIndex = 0;
count++;
// Wake a reader after insertion.
notEmpty.signal();
}
3.4. take [Blocking]
public E take() throws InterruptedException {
// Lock.
final ReentrantLock lock = this.lock;
lock.lockInterruptibly();
try {
// If the array is empty, wait. A writer wakes the reader after adding an element.
while (count == 0)
notEmpty.await();
// Dequeue.
return dequeue();
} finally {
// Unlock.
lock.unlock();
}
}
- Line 3: acquire the lock.
- Lines 6-8: if the array is empty, wait.
- Line 10: if it is not empty, dequeue and wake a writer.
3.4.1. Lock
final ReentrantLock lock = this.lock;
lock.lockInterruptibly();
try {
//...
} finally {
lock.unlock();
}
3.4.2. If the Array Is Empty, Wait
while (count == 0)
notEmpty.await();
3.4.3. If It Is Not Empty, Dequeue and Wake Writers
- dequeue
private E dequeue() {
final Object[] items = this.items;
@SuppressWarnings("unchecked")
// Obtain the next element and set the array slot to null.
E x = (E) items[takeIndex];
items[takeIndex] = null;
// After reaching the end, reset takeIndex to 0.
// This array is reused circularly and does not need expansion.
if (++takeIndex == items.length)
takeIndex = 0;
count--;
if (itrs != null)
itrs.elementDequeued();
// Wake a writer after dequeue.
notFull.signal();
return x;
}
3.5. offer [Returns a Special Value]
public boolean offer(E e) {
checkNotNull(e);
// Lock.
final ReentrantLock lock = this.lock;
lock.lock();
try {
// Full: return false directly.
if (count == items.length)
return false;
else {
// Not full: enqueue and wake a reader.
enqueue(e);
return true;
}
} finally {
// Unlock.
lock.unlock();
}
}
3.6. poll [Returns a Special Value]
public E poll() {
final ReentrantLock lock = this.lock;
// Lock.
lock.lock();
try {
// Return null if empty; otherwise dequeue and wake a writer.
return (count == 0) ? null : dequeue();
} finally {
lock.unlock();
}
}
3.7. add [Throws an Exception]
public boolean add(E e) {
// Simply call AbstractQueue.add.
return super.add(e);
}
// AbstractQueue.add
public boolean add(E e) {
// Call ArrayBlockingQueue.offer.
if (offer(e))
return true;
else
throw new IllegalStateException("Queue full");
}
3.8. remove [Throws an Exception]
public E remove() {
// Simply call poll.
E x = poll();
if (x != null)
return x;
else
// If there is no element, throw an exception.
throw new NoSuchElementException();
}
3.9. element [Throws an Exception]
public E element() {
// Call peek.
E x = peek();
if (x != null)
return x;
else
// If empty, throw an exception.
throw new NoSuchElementException();
}
3.10. peek [Returns a Special Value]
public E peek() {
// Lock.
final ReentrantLock lock = this.lock;
lock.lock();
try {
return itemAt(takeIndex); // null when queue is empty
} finally {
// Unlock.
lock.unlock();
}
}
@SuppressWarnings("unchecked")
final E itemAt(int i) {
// Return the i-th element directly.
return (E) items[i];
}
4. Summary
It is implemented using an array and is a bounded queue.
It uses one lock and two conditions. One lock means producers and consumers share mutual exclusion while modifying queue state. The two conditions allow readers and writers to wake one another according to whether the queue is empty or full.
Discussion
Sign in with GitHub to comment. Discussions are stored as GitHub Issues.View on GitHub