NOTE

6.35 LinkedBlockingQueue

What LinkedBlockingQueue is, usage, and source analysis of its linked-list structure, separate put/take locks, conditions, and queue operations.

JavaCreated Updated 2 min readhistorical

This is a historical learning note and may contain outdated or incomplete understanding.

1. What Is It?

A blocking queue implemented with a singly linked list.

It can be bounded by passing a capacity; the no-argument constructor uses Integer.MAX_VALUE as capacity.

Reads compete with reads through the take lock, writes compete with writes through the put lock, while a reader and writer can operate concurrently because separate locks are used.

Its throughput can therefore be higher than ArrayBlockingQueue in some producer/consumer workloads.

2. How to Use It

public class LinkedBlockingQueueTest
{
    public static void main(String[] args) throws InterruptedException
    {
        LinkedBlockingQueue<String> queue = new LinkedBlockingQueue<>(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. Source-Code Analysis

3.1. Constructor

3.1.1. Implemented with Singly Linked List + Lock + Condition

public class LinkedBlockingQueue<E> extends AbstractQueue<E>
        implements BlockingQueue<E>, java.io.Serializable {

    // Maximum length.
    private final int capacity;

    // Actual length.
    private final AtomicInteger count = new AtomicInteger();

    // Head node.
    transient Node<E> head;

    // Tail node.
    private transient Node<E> last;

    // Lock used for dequeue operations; protects the head.
    private final ReentrantLock takeLock = new ReentrantLock();

    // If the queue is empty during a read, wait on notEmpty.
    private final Condition notEmpty = takeLock.newCondition();

    // Lock used for enqueue operations; protects the tail.
    private final ReentrantLock putLock = new ReentrantLock();

    // If the queue is full during a write, wait on notFull.
    private final Condition notFull = putLock.newCondition();

    public LinkedBlockingQueue() {
        // Effectively a very large bounded queue.
        this(Integer.MAX_VALUE);
    }

    public LinkedBlockingQueue(int capacity) {
        if (capacity <= 0) throw new IllegalArgumentException();
        this.capacity = capacity;//Bounded queue.
        last = head = new Node<E>(null);//Head is a sentinel node.
    }
}

3.1.2. Node

static class Node<E> {
    E item;

    // Singly linked queue.
    Node<E> next;

    Node(E x) { item = x; }
}

The structure is shown below:

3.2. put [Blocking]

public void put(E e) throws InterruptedException {
    if (e == null) throw new NullPointerException();
    int c = -1;
    Node<E> node = new Node<E>(e);
    final ReentrantLock putLock = this.putLock;
    final AtomicInteger count = this.count;
    // Acquire write lock.
    putLock.lockInterruptibly();
    try {
        // If current size reaches capacity, block until a reader takes an element.
        while (count.get() == capacity) {
            notFull.await();
        }
        // Append at tail.
        enqueue(node);
        c = count.getAndIncrement();//Increment, but return old value.
        if (c + 1 < capacity)
            notFull.signal();//Wake another writer.
    } finally {
        putLock.unlock();
    }
    // c == 0 means the queue was previously empty, so readers may be waiting.
    if (c == 0)
        // Wake a thread waiting in poll/take.
        signalNotEmpty();
}
  • Line 8: acquire the write lock. Other writers cannot enter simultaneously, but readers using takeLock can still run.
  • Lines 11-13: if the queue is full, the writer blocks waiting for a reader.
  • Line 15: if not full, append the element at the tail.
  • Line 16: increment count and obtain the old value.
  • Lines 17-18: if the queue is still not full, wake another writer.
  • Lines 23-25: if the old count was 0, the queue was previously empty, so wake a waiting reader.

3.2.1. Acquire the Write Lock

putLock.lockInterruptibly();
try {
    //...
} finally {
    putLock.unlock();
}

3.2.2. If the Queue Is Full, Wait

while (count.get() == capacity) {
    notFull.await();
}

3.2.3. If Not Full, Enqueue

  • enqueue
private void enqueue(Node<E> node) {
    // Append the node at the tail and update last.
    last = last.next = node;
}

3.2.4. If the Queue Is Still Not Full After Enqueue, Wake Another Writer

if (c + 1 < capacity)
    notFull.signal();

3.2.5. If the Queue Was Previously Empty, Wake a Reader After Unlocking

if (c == 0)
    signalNotEmpty();
  • signalNotEmpty
private void signalNotEmpty() {
    // Acquire read lock.
    final ReentrantLock takeLock = this.takeLock;
    takeLock.lock();
    try {
        // Wake a reader.
        notEmpty.signal();
    } finally {
        takeLock.unlock();
    }
}

3.3. take [Blocking]

public E take() throws InterruptedException {
    E x;
    int c = -1;
    final AtomicInteger count = this.count;
    final ReentrantLock takeLock = this.takeLock;
    // Acquire read lock.
    takeLock.lockInterruptibly();
    try {
        // If size is 0, block until a writer inserts an element.
        while (count.get() == 0) {
            notEmpty.await();
        }
        // Remove first real node.
        x = dequeue();
        c = count.getAndDecrement();
        if (c > 1)
            notEmpty.signal();//Wake another reader.
    } finally {
        takeLock.unlock();
    }
    // c == capacity means the queue was previously full, so a writer may be waiting.
    if (c == capacity)
        // Wake a thread waiting in put.
        signalNotFull();
    return x;
}
  • Line 7: acquire the read lock. Other readers cannot enter simultaneously, but writers using putLock can still run.
  • Lines 10-12: if the queue is empty, the reader waits for a writer.
  • Line 14: if not empty, remove the head element.
  • Line 15: decrement count and return the old value.
  • Lines 16-17: if the queue is still non-empty, wake another reader.
  • Lines 22-24: if the old count equaled capacity, the queue had been full, so wake a waiting writer.

3.3.1. Acquire the Read Lock

takeLock.lockInterruptibly();
try {
    //....
} finally {
    takeLock.unlock();
}

3.3.2. If the Queue Is Empty, Wait

while (count.get() == 0) {
    notEmpty.await();
}

3.3.3. If Not Empty, Dequeue

  • dequeue
private E dequeue() {
    Node<E> h = head;//Head is a sentinel node.
    Node<E> first = h.next;//The actual first node.
    h.next = h; // help GC
    head = first;//Move head to the first real node.
    E x = first.item;
    first.item = null;
    return x;
}

3.3.4. If the Queue Is Still Non-Empty, Wake Another Reader

if (c > 1)
    notEmpty.signal();

3.3.5. If the Queue Was Previously Full, Wake a Writer After Unlocking

if (c == capacity)
    signalNotFull();
  • signalNotFull
private void signalNotFull() {
    final ReentrantLock putLock = this.putLock;
    // Acquire write lock.
    putLock.lock();
    try {
        // Notify a writer that space is available.
        notFull.signal();
    } finally {
        putLock.unlock();
    }
}

3.4. offer [Returns a Special Value]

public boolean offer(E e) {
    if (e == null) throw new NullPointerException();
    final AtomicInteger count = this.count;
    if (count.get() == capacity)
        return false;
    int c = -1;
    Node<E> node = new Node<E>(e);
    final ReentrantLock putLock = this.putLock;
    putLock.lock();
    try {
        if (count.get() < capacity) {
            enqueue(node);
            c = count.getAndIncrement();
            if (c + 1 < capacity)
                notFull.signal();
        }
    } finally {
        putLock.unlock();
    }
    if (c == 0)
        signalNotEmpty();
    return c >= 0;//Difference from put: return instead of blocking.
}

3.5. poll [Returns a Special Value]

public E poll() {
    final AtomicInteger count = this.count;
    if (count.get() == 0)//Difference from take: return null.
        return null;
    E x = null;
    int c = -1;
    final ReentrantLock takeLock = this.takeLock;
    takeLock.lock();
    try {
        if (count.get() > 0) {
            x = dequeue();
            c = count.getAndDecrement();
            if (c > 1)
                notEmpty.signal();
        }
    } finally {
        takeLock.unlock();
    }
    if (c == capacity)
        signalNotFull();
    return x;
}

3.6. peek [Returns a Special Value]

public E peek() {
    if (count.get() == 0)//Return null if empty.
        return null;
    final ReentrantLock takeLock = this.takeLock;
    takeLock.lock();
    try {
        Node<E> first = head.next;
        if (first == null)
            return null;
        else
            return first.item;
    } finally {
        takeLock.unlock();
    }
    // No need to wake a writer because no element was dequeued.
}

4. Summary

It is implemented with a singly linked list and can be configured with a bounded capacity; the no-argument constructor uses a very large capacity.

It uses two locks and two conditions. The separate read/write locks allow a reader and writer to proceed concurrently, while the two conditions let readers and writers wake one another when queue state changes.

5. References

Discussion

Sign in with GitHub to comment. Discussions are stored as GitHub Issues.View on GitHub