SynchronousQueue / LinkedTransferQueue 无锁队列源码
概述
SynchronousQueue 与 LinkedTransferQueue 是 JUC 中两员"直接交接"的队列:
SynchronousQueue:0 容量队列——put必须等到take同时到达,元素不落地,直接在线程间交接。LinkedTransferQueue:无锁无界队列的"传输版"——transfer优先与等待的消费者直接交接,没有消费者时才入队挂起。
两者都基于"节点模式 + CAS 配对/追加"的无锁算法(Doug Lea 设计),没有锁、没有 Condition,只有自旋与 LockSupport.park。本文基于 OpenJDK 21 源码拆解这两种队列的核心实现。
核心源码解析
① SynchronousQueue 的 TransferStack 与 TransferQueue
java
public SynchronousQueue() { this(false); } // 默认非公平
public SynchronousQueue(boolean fair) {
transferer = fair ? new TransferQueue<E>() : new TransferStack<E>();
}| 实现 | 模式 | 结构 | 语义 |
|---|---|---|---|
TransferStack | 非公平(默认) | 栈(LIFO) | 后到的请求优先配对 |
TransferQueue | 公平 | 队列(FIFO) | 先到的请求优先配对 |
Transferer<E>是公共抽象:abstract E transfer(E e, boolean timed, long nanos),e == null表示 take(REQUEST),否则表示 put(DATA)。- LIFO vs FIFO 的取舍:栈模式吞吐更高(热数据在栈顶,缓存友好)但存在"插队"不公平;队列模式严格公平但争用略高。
transfer是所有操作(put/take/offer/poll)的公共入口,timed/nanos 区分阻塞与超时。
② TransferStack.transfer(SNode s, boolean timed, long nanos) 的三状态
java
static final class TransferStack<E> extends Transferer<E> {
static final int REQUEST = 0; // 消费者(take)
static final int DATA = 1; // 生产者(put)
static final int FULFILLING = 2; // 配对进行中
volatile SNode head; // 栈顶
E transfer(E e, boolean timed, long nanos) {
SNode s = null;
int mode = (e == null) ? REQUEST : DATA;
for (;;) {
SNode h = head;
if (h == null || h.mode == mode) { // ① 栈空或同模式:入栈等待
if (s == null) s = new SNode(e);
if (casHead(h, s)) {
SNode m = awaitFulfill(s, timed, nanos); // 等待配对
if (m == s) { clean(s); return null; } // 超时/中断:清理并返回
... // 出栈
}
} else if (!isFulfilling(h.mode)) { // ② 反模式且未配对:尝试配对
...
if (h.next == s && casHead(h, s.tryMatch(h))) { ... } // 压入 FULFILLING 节点
} else { // ③ 头节点正在配对:帮助完成
SNode m = h.next;
if (m == null) casHead(h, null);
else if (m.match == h) casHead(h, m.next); // 帮配对节点出栈
else m.tryMatch(h);
}
}
}
}- 三种角色:
REQUEST(等待取)、DATA(等待放)、FULFILLING(与某个等待者配对中)。配对完成的标志是节点的match字段指向对方。 - 配对流程:新请求与栈顶模式相反时,压入一个
FULFILLING节点,tryMatch把双方match互相指向并 CAS 出栈——一次 CAS 完成"成交"。 - 帮助机制:任何线程看到
FULFILLING头节点都会帮它完成出栈,保证配对不被卡住(无锁队列的典型"协作"设计)。
③ TransferStack.awaitFulfill(SNode s, boolean timed, long nanos) 的自旋等待
java
E awaitFulfill(SNode s, boolean timed, long nanos) {
final long deadline = timed ? System.nanoTime() + nanos : 0L;
Thread w = Thread.currentThread();
int spins = (shouldSpin(s) ? (timed ? MAX_TIMED_SPINS : MAX_UNTIMED_SPINS) : 0);
for (;;) {
if (w.isInterrupted()) s.tryCancel(); // 中断 → 尝试取消
SNode m = s.match;
if (m != null) return m; // 配对成功
if (timed) { ... } // 超时判断
if (spins > 0) { --spins; Thread.onSpinWait(); } // ① 自旋
else if (s.waiter == null) s.waiter = w; // ② 登记等待线程
else if (!timed) LockSupport.park(this); // ③ 无限阻塞
else if (nanos > spinForTimeoutThreshold) LockSupport.parkNanos(this, nanos);
}
}- 自旋阈值
spinForTimeoutThreshold(约 1000ns):等待时间短于阈值时纯自旋,避免 park/unpark 的系统调用开销。 waiter字段登记等待线程,配对成功方通过LockSupport.unpark精确唤醒。match是非空即成功:被配对时match指向对方节点,自旋循环立即返回。
④ SynchronousQueue 的 0 容量语义
- 无内部存储:任何时刻队列中至多只有"等待配对"的线程,元素不进入任何容器——
take()必须等到put()到达,反之亦然。 - 正因为 0 容量,
offer(e)(NOW 模式)在无等待消费者时立即返回 false;put(e)(SYNC 模式)则一直阻塞到成功交接。 - 适用场景:线程池
Executors.newCachedThreadPool()的SynchronousQueue任务队列——提交任务时若无空闲工作线程,直接失败并触发线程扩容;交接零拷贝、零延迟。 - 与
TransferStack的"配对即完成"不同,队列模式的TransferQueue用QNode的isData标记区分生产/消费节点,FIFO 顺序配对。
⑤ LinkedTransferQueue.xfer(E e, boolean haveData, int how, long nanos) 的通用方法
java
private static final int NOW = 0; // 立即:无消费者则返回
private static final int ASYNC = 1; // 异步:无消费者则入队
private static final int SYNC = 2; // 同步:等待交接
private static final int TIMED = 3; // 超时同步
E xfer(E e, boolean haveData, int how, long nanos) {
Node s = null;
for (;;) {
for (Node h = head, p = h; p != null; ) { // ① 尝试匹配队首反模式节点
boolean isData = p.isData;
Object item = p.item;
if (item != p && (item != null) == isData) { // 未匹配节点
if (isData == haveData) break; // 模式相同:停止扫描
if (p.casItem(item, e)) { // CAS 完成交换
... // 更新 head
if (p.waiter != null) LockSupport.unpark(p.waiter); // 唤醒等待者
return (E) item;
}
}
Node n = p.next;
p = (p != n) ? n : (h = head); // 跳过已删除节点
}
if (how != NOW) { // ② 无匹配:追加节点
if (s == null) s = new Node(e, haveData);
Node pred = tryAppend(s, haveData);
if (pred != null) {
if (how == ASYNC) return e; // 入队即返回
return awaitMatch(s, pred, e, (how == TIMED), nanos); // ③ 同步等待
}
}
return e; // NOW 且无法匹配
}
}- 四模式入口:
transfer→SYNC、tryTransfer→NOW、tryTransfer(timeout)→TIMED、put→SYNC、offer→NOW/ASYNC、poll/take→NOW/SYNC——一个xfer撑起整个 API。 - 扫描匹配:从 head 沿链表找反模式节点,
casItem(item, e)交换值完成交接;sweep()/unsplice()负责惰性清理已删除节点(超过SWEEP_THRESHOLD触发全链清扫)。 - 返回值语义:
haveData(put)成功返回e;take 成功返回取到的值。
⑥ LinkedTransferQueue.tryAppend(SNode s, boolean haveData) 的节点追加
java
private Node tryAppend(Node s, boolean haveData) {
for (Node t = tail, p = t;;) {
Node n;
if (p == null && (p = head) == null) { // 空队列:s 成为唯一节点
if (casHead(null, s)) return s;
} else if (p.cannotPrecede(haveData)) { // 与 head 反模式冲突:取消尾链接
...
} else if ((n = p.next) != null) { // 尾部已推进:跟到新尾
p = p != t && t != (u = tail) ? (t = u) : (p != n) ? n : (h = head);
} else if (!p.casNext(null, s)) { // CAS 追加失败:重试
...
} else {
if (p != t) casTail(t, s); // 追加成功:更新 tail
return s; // 返回前驱
}
}
}- CAS 追加到尾部:
p.casNext(null, s)把新节点链到末节点,随后casTail更新尾指针——无锁队列的经典两阶段追加。 cannotPrecede检查:若 head 是反模式(如全是消费者却来了消费者),说明链表顶部存在失效节点,需要先清理再追加。- 返回值是前驱节点:
awaitMatch用它做"pred 被删除"的感知(sweep时从 pred 处重链)。
⑦ LinkedTransferQueue.awaitMatch(SNode s, boolean timed, long nanos) 的阻塞
java
private E awaitMatch(Node s, Node pred, E e, boolean timed, long nanos) {
final long deadline = timed ? System.nanoTime() + nanos : 0L;
Thread w = Thread.currentThread();
int spins = -1;
for (;;) {
Object item = s.item;
if (item != e) { // ① 已被匹配:item 被交换
s.waiter = null;
return (E) item;
}
if (w.isInterrupted()) { ... } // 中断处理
if (spins < 0) spins = ...; // 初始化自旋次数
else if (spins > 0) { --spins; Thread.onSpinWait(); } // ② 自旋
else if (s.waiter == null) s.waiter = w; // ③ 登记等待线程
else if (!timed) LockSupport.park(this); // ④ 阻塞
else if (nanos > spinForTimeoutThreshold) LockSupport.parkNanos(this, nanos);
}
}- 唤醒条件
s.item != e:匹配方通过casItem(e, 对方值)改变节点的 item,等待方自旋循环观察到 item 变化即返回——值变化本身就是唤醒信号。 Thread.onSpinWait()是 x86pause指令的 Java 暴露,告诉 CPU"正在自旋",避免忙等拖垮流水线。- 与
awaitFulfill同理:短等待自旋、长等待 park,配合waiter登记实现精准unpark。
⑧ LinkedTransferQueue.getWaitingConsumerCount() / hasWaitingConsumer()
java
public int getWaitingConsumerCount() {
int count = 0;
for (Node p = head; p != null; p = p.next) {
if (p.isData == false && p.waiter != null) // 消费者请求节点 + 已登记等待线程
++count;
}
return count;
}
public boolean hasWaitingConsumer() {
return getWaitingConsumerCount() != 0;
}- 判断依据:请求节点
isData == false且waiter != null——即"等待中的消费者"。生产者可用它决定"直接 transfer(有消费者)"还是"入队缓冲(无消费者)"。 - 这是
LinkedTransferQueue独有的探测能力,SynchronousQueue不提供(其语义不允许"有容量"判断)。 - 注意计数是遍历快照,非精确保证;
hasWaitingConsumer常用于动态选择传输策略(如无消费者时用批量队列,有消费者时走 transfer 直达)。
总结
| 队列 | 存储 | 交接方式 | 关键类 |
|---|---|---|---|
SynchronousQueue(栈) | 0 容量 | LIFO 栈 + FULFILLING 配对 | TransferStack / SNode |
SynchronousQueue(公平) | 0 容量 | FIFO 队 + QNode 匹配 | TransferQueue / QNode |
LinkedTransferQueue | 无界 | 优先匹配,否则入队等待 | Node + xfer 四模式 |
三者的共性是无锁设计精髓:节点即消息(item 的 CAS 交换同时完成交接与唤醒)、协作推进(帮助他人完成配对)、自旋降级(短等待自旋、长等待 park)。SynchronousQueue 用于零缓冲交接,LinkedTransferQueue 在 transfer 语义下把"直达 + 缓冲"两种模式统一到一个 xfer 中,是 JUC 无锁算法的高阶范例。