NOTE

6.42 SynchronousQueue

1. What it is. A blocking queue implemented internally with a singly linked structure that does not store elements. A writer must have a reader at the same time in order to proceed, and vice versa; otherwise the writer or reader remains blocked. 2. Usage. 3. Principles. 3.1. Constructor. 3.1.1. Transfer. 3.1.2. QNode. 3.2. put blocking. 3.2.1. Calls TransferQueue.

JavaCreated Updated 1 min readhistorical

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

1. What It Is

A blocking queue implemented internally with a singly linked structure. It does not store elements.

A writer must have a reader at the same time in order to proceed, and vice versa.

Otherwise, the writer will remain blocked, or the reader will remain blocked.

2. Usage

public class SynchronousQueueTest
{
    public static void main(String[] args) throws InterruptedException
    {
        SynchronousQueue <String> queue = new SynchronousQueue<>();
        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);//peek is not supported
                    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. Principles

3.1. Constructor

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

    private transient volatile Transferer<E> transferer;

    public SynchronousQueue() {
        this(false);// Nonfair by default, which means using a stack
    }

    public SynchronousQueue(boolean fair) {
        transferer = fair ? new TransferQueue<E>() : new TransferStack<E>();
    }

    // Head and tail of the singly linked list
    transient volatile QNode head;
    transient volatile QNode tail;
}

3.1.1. Transfer

abstract static class Transferer<E> {
     // Both put and take call this method
     // If e is null, this represents a reader's take operation
     // If e is not null, this represents a writer's put operation
     // The second parameter indicates whether a timeout is set.
     // If a timeout is set, the third parameter is the timeout value.
     // If the return value is null, it means either timeout or interruption.
     // Which one occurred can be determined by checking the interrupt status.
    abstract E transfer(E e, boolean timed, long nanos);
}

3.1.2. QNode

static final class QNode {
    volatile QNode next;          // Singly linked list
    volatile Object item;         // CAS'ed to or from null
    volatile Thread waiter;       // to control park/unpark
    final boolean isData;// true means write, false means read
}

3.2. put Blocks

public void put(E e) throws InterruptedException {
    // The writer guarantees that e is not null
    if (e == null) throw new NullPointerException();
    // Call Transfer.transfer to pass the element to a reader
    if (transferer.transfer(e, false, 0) == null) {
        Thread.interrupted();
        throw new InterruptedException();
    }
}

3.2.1. Calls TransferQueue

  • TransferQueue.transfer
E transfer(E e, boolean timed, long nanos) {
    QNode s = null;
    boolean isData = (e != null);// e is non-null for a write (true), null for a read (false)

    for (;;) {
        QNode t = tail;
        QNode h = head;
        if (t == null || h == null)         // saw uninitialized value
            continue;                       // spin

        // The queue is empty, or the mode of the tail node is the same as the current node
        // (that is, both are writes or both are reads).
        // In that case, directly enqueue the current node.
        if (h == t || t.isData == isData) {
            QNode tn = t.next;
            // The previous tail differs from the current tail, which means another node has already been enqueued.
            // Start over.
            if (t != tail)
                continue;
            // At this point tail has not changed, but tail.next is non-null.
            // This means a node has been enqueued but tail has not yet been updated.
            // Point tail to tail.next.
            if (tn != null) {
                advanceTail(t, tn);// If tail == t, point tail to tn
                continue;
            }
            // A timeout was configured, but no time remains
            if (timed && nanos <= 0)
                return null;
            // Construct the current node
            if (s == null)
                s = new QNode(e, isData);
            // Insert it at the end of the linked list
            if (!t.casNext(null, s))
                continue;

            // If tail == t, point tail to s
            advanceTail(t, s);
            // Spin or block waiting for a thread of the opposite mode to arrive and wake this one
            // A writer gets null; a reader gets the writer's value
            Object x = awaitFulfill(s, e, timed, nanos);

            // Reaching here means the thread has been awakened; continue below
            if (x == s) {                   // wait was cancelled
                clean(t, s);
                return null;
            }

            // The current node's next is not the current node
            // When the head node == tail node, CAS the head to the current node
            if (!s.isOffList()) {           // not already unlinked
                advanceHead(t, s);          // unlink if head
                if (x != null)              // and forget fields
                    s.item = s;
                s.waiter = null;
            }
            return (x != null) ? (E)x : e;

        }
        // One read and one write match exactly
        else {
            // The node after head is the current node
            QNode m = h.next;
            // If head or tail changed, or head.next became null, the linked list changed.
            // Start over.
            if (t != tail || m == null || h != head)
                continue;

            // Retry on failure
            Object x = m.item;
            if (isData == (x != null) ||    // m already fulfilled
                x == m ||                   // m cancelled
                !m.casItem(x, e)) {         // lost CAS
                advanceHead(h, m);          // dequeue and retry
                continue;
            }

            // Change the head node with CAS. If h == head, set head to the current node
            advanceHead(h, m);              // successfully fulfilled
            // Wake the thread associated with the current node. Corresponds to awaitFulfill
            LockSupport.unpark(m.waiter);

            return (x != null) ? (E)x : e;
        }
    }
}
  • advanceTail
void advanceTail(QNode t, QNode nt) {
    // If the current tail node == the passed-in tail node
    if (tail == t)
        // Use CAS to change the tail pointer to nt
        UNSAFE.compareAndSwapObject(this, tailOffset, t, nt);
}
  • awaitFulfill
// Either spin or block
Object awaitFulfill(QNode s, E e, boolean timed, long nanos) {
    // If a timeout is configured, calculate the deadline
    final long deadline = timed ? System.nanoTime() + nanos : 0L;
    Thread w = Thread.currentThread();
    // If head.next is the current node, do not block immediately; spin and wait
    int spins = ((head.next == s) ?
                 (timed ? maxTimedSpins : maxUntimedSpins) : 0);
    for (;;) {
        // If the current thread has been interrupted, CAS the current node's item field to e
        if (w.isInterrupted())
            s.tryCancel(e);
        // This is the only exit from this method
        // Return when the current node's item field differs from e
        Object x = s.item;
        if (x != e)
            return x;
        // If timed, update the remaining timeout
        if (timed) {
            nanos = deadline - System.nanoTime();
            if (nanos <= 0L) {
                s.tryCancel(e);
                continue;
            }
        }
        // Decrease the spin count by 1 each loop
        if (spins > 0)
            --spins;
        // Reaching here means the maximum spin count has been exhausted
        // or spinning was not configured

        // If the current node has not yet been associated with a thread, associate it
        else if (s.waiter == null)
            s.waiter = w;
        // If no timeout is configured, block
        else if (!timed)
            LockSupport.park(this);
        else if (nanos > spinForTimeoutThreshold)
            LockSupport.parkNanos(this, nanos);
    }
}
  • tryCancel
void tryCancel(Object cmp) {
    UNSAFE.compareAndSwapObject(this, itemOffset, cmp, this);
}
  • advanceHead
void advanceHead(QNode h, QNode nh) {
    if (h == head &&
        UNSAFE.compareAndSwapObject(this, headOffset, h, nh))
        h.next = h; // forget old next
}

3.3. take Blocks

public E take() throws InterruptedException {
    // Call Transfer.transfer to obtain an element from a writer
    E e = transferer.transfer(null, false, 0);
    if (e != null)
        return e;
    Thread.interrupted();
    throw new InterruptedException();
}

4. Summary

It does not store elements, and its throughput is higher than LinkedBlockingQueue.

Reads and writes must match before either side can proceed. Otherwise, the operation is added to the queue and blocks until a thread of the opposite mode arrives and wakes it.

5. References

Discussion

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