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.
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.
Discussion
Sign in with GitHub to comment. Discussions are stored as GitHub Issues.View on GitHub