本文默认读者已掌握线程生命周期、JMM 与 synchronized/volatile(系列第十二篇),JDK 以 8 为主,涉及 9 之后的变化(Unsafe 迁往 VarHandle)会单独标注。AQS 一节最硬,建议按”state 是什么 → 队列怎么排 → 线程怎么睡 → 谁把它叫醒”四步推演,读完再回头看 ReentrantLock、CountDownLatch、Semaphore,会发现它们只是同一套骨架上的不同血肉。
一、JUC 全景与定位 靠 synchronized 与 volatile 能解决大部分并发问题,但它们有两个天然短板:一是不可中断 ——等监视器锁时无法响应 interrupt(),死锁只能重启;二是不够灵活 ——没有尝试加锁、没有超时、没有读写分离、也没有多个等待队列,想做”生产者只唤醒消费者”这种精细控制,wait/notify 只能 notifyAll 全量广播。
JUC(java.util.concurrent)正是为了补上这块拼图。它由 Doug Lea 主导设计,从 JDK 5 引入,核心思路是:把并发控制的公共骨架抽出来做成可复用组件 ,让业务代码从”手工操作 wait/notify”升级到”组合现成工具”。
1.1 包的四层结构
子包/包
代表类
解决的核心问题
本篇覆盖
java.util.concurrent.locks
AbstractQueuedSynchronizer、ReentrantLock、ReentrantReadWriteLock、StampedLock
显式锁、条件队列、读写分离、乐观读
二、三章
java.util.concurrent.atomic
AtomicInteger、AtomicReference、LongAdder、AtomicStampedReference
无锁原子更新、累加器
四章
java.util.concurrent(集合部分)
ConcurrentHashMap、CopyOnWriteArrayList、ConcurrentLinkedQueue、各类 BlockingQueue
线程安全容器的性能与语义
五章
java.util.concurrent(协作部分)
CountDownLatch、CyclicBarrier、Semaphore、Phaser、Exchanger
线程间的等待、汇合、限流
六章
java.util.concurrent(执行器)
ThreadPoolExecutor、FutureTask、CompletableFuture、ForkJoinPool
任务调度与异步编排
下一篇
附带
ThreadLocal、ThreadLocalRandom
线程封闭、上下文传递
七章
1.2 与 synchronized 的分工 很多人纠结”有了 synchronized 为什么还要 Lock”。答案不是替换,而是分工 :synchronized 是 JVM 内置的 Monitor 机制,解锁由编译器插入的 monitorenter/monitorexit 保证,JVM 还会做锁粗化、锁消除与自适应自旋,JDK 6 之后性能已持平,能用就用 ;Lock 是纯 Java 实现的显式锁,胜在能力——可中断、可超时、可非阻塞尝试、可选公平性、一把锁绑多个 Condition、支持读写分离与乐观读。需要其中一项时再切换。
选择顺序建议:先想能不能”不加锁”(无锁结构、线程封闭、不可变对象)→ 再想 synchronized 够不够 → 需要高级能力才上 ReentrantLock → 读多写极少考虑 StampedLock 或 CopyOnWrite → 高并发计数优先 LongAdder。
1.3 本篇与下一篇的边界
关注点
本篇(十三)
下一篇(十四)
核心主题
同步原语与容器:如何”安全地共享数据”
任务执行框架:如何”高效地执行任务”
关键抽象
AQS、Lock、Condition、原子类、BlockingQueue
Executor、Future、CompletableFuture、ForkJoinPool
典型问题
竞态、可见性、死锁、伪共享、内存泄漏
线程池参数、拒绝策略、异步编排、任务窃取
交集
BlockingQueue 既是容器也是线程池的工作队列;FutureTask 内部也用到了 AQS 的共享模式,本篇会顺带点出
线程池的 worker 抢任务机制会复用本篇的 CAS 知识
二、AQS 原理:一把锁的骨架 AbstractQueuedSynchronizer(下文简称 AQS)是 JUC 的心脏 。ReentrantLock、Semaphore、CountDownLatch、ReentrantReadWriteLock、FutureTask、ThreadPoolExecutor.Worker 全部直接或间接继承自它。理解 AQS,等于一次性理解了半个 JUC。
2.1 模板方法模式 AQS 用了非常经典的模板方法模式 :把排队、阻塞、唤醒、取消这些与业务无关的流程全部实现在基类,只把”能不能拿到资源”这一个判断留给子类。子类需重写的方法只有五个,全是 protected:
方法
模式
语义
tryAcquire(int)
独占
尝试获取资源,成功返回 true
tryRelease(int)
独占
尝试释放资源
tryAcquireShared(int)
共享
返回负数失败,0 成功但无剩余,正数成功且有剩余
tryReleaseShared(int)
共享
释放共享资源
isHeldExclusively()
独占
当前线程是否独占持有,供 Condition 使用
其余 acquire、acquireInterruptibly、acquireShared、release 等全是 final 模板方法,子类不可改。这就是 AQS 的优雅之处:流程固定,语义可插拔 。
2.2 核心三件套 2.2.1 volatile int state state 是同步器的资源计数器 ,volatile 保证可见性,配 getState/setState/compareAndSetState 访问(CAS 由 Unsafe 提供)。其含义由子类定义:
ReentrantLock:state=0 表示未锁定,>0 表示被持有,数值即重入次数;
Semaphore:state 表示剩余许可数;
CountDownLatch:state 表示还没完成的计数;
ReentrantReadWriteLock:state 高 16 位存读锁数量,低 16 位存写锁重入次数——一个 int 拆成两个半字用。
2.2.2 CLH 队列的变体 AQS 内部维护一条 FIFO 双向队列 ,官方注释称之为 “variant of CLH queue”。相对原始 CLH(单向、靠忙等前驱状态位)做了三点改造:
改双向 :Node 同时持有 prev 与 next。单向链表在节点 CANCELLED 时拿不到前驱,无法从中间摘除;双向可直接 prev.next = next; next.prev = prev。
自旋改 park :竞争激烈、临界区长时忙等纯属烧 CPU,AQS 改为自旋失败后 LockSupport.park() 挂起,由前驱释放时 unpark。
引入 head 哑节点 :head 指向”当前已持有资源的线程”或空壳节点,真正排队的是 head.next。
不用普通 LinkedList 的原因是无锁入队:enq 用 CAS 抢设 tail,失败自旋重试,全程不加锁;而 prev 指针让”前驱取消就往前跳”实现得非常干净。
2.2.3 Node 与 waitStatus 1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 static final class Node { static final Node SHARED = new Node (); static final Node EXCLUSIVE = null ; static final int CANCELLED = 1 ; static final int SIGNAL = -1 ; static final int CONDITION = -2 ; static final int PROPAGATE = -3 ; volatile int waitStatus; volatile Node prev; volatile Node next; volatile Thread thread; Node nextWaiter; }
waitStatus 是理解 AQS 的钥匙,五个取值必须记牢:
CANCELLED = 1 :唯一正数。节点因超时或中断放弃竞争,一旦置位不再变化,会在 cancelAcquire 中被摘链。
SIGNAL = -1 :最核心 。表示”我的后继在 park,我有义务释放时 unpark 它”。所以每个节点入队后必须先确保前驱是 SIGNAL 才敢 park,否则会”睡着了没人叫”。
CONDITION = -2 :节点不在同步队列,而在某个 ConditionObject 的单向等待队列里,signal 时搬回同步队列并把状态改回 0。
PROPAGATE = -3 :共享模式专用。多个线程可同时持有资源,一次 releaseShared 可能要连续唤醒多个后继,该状态保证唤醒能向后传播,避免丢失。
0 :新建节点的默认值,表示”当前无事发生”。
2.3 acquire 全流程推演 独占模式获取资源的入口是 acquire(int arg),逻辑只有三行,但背后是完整的一套排队机制:
1 2 3 4 5 public final void acquire (int arg) { if (!tryAcquire(arg) && acquireQueued(addWaiter(Node.EXCLUSIVE), arg)) selfInterrupt(); }
完整流程分五步:
快速尝试 :调子类 tryAcquire。非公平锁下这步常直接成功(刚释放的锁被新线程抢走),这是非公平吞吐更高的根源。
包装入队 :addWaiter 把当前线程包成 EXCLUSIVE 节点,先 CAS 一次挂队尾,失败则进 enq 自旋 CAS 直到成功。
排队自旋 :acquireQueued 是 for(;;):前驱是 head 就再 tryAcquire,失败则问 shouldParkAfterFailedAcquire 能否 park。
安全入睡 :shouldParkAfterFailedAcquire 先把前驱 waitStatus CAS 成 SIGNAL 并返回 false 让外层再转一圈;第二圈发现前驱已是 SIGNAL,才真正 LockSupport.park(this)。
被唤醒重来 :前驱释放时 unparkSuccessor 找到 head 后第一个非 CANCELLED 节点并 unpark。被唤醒线程回到第 3 步,tryAcquire 成功后自己成为新 head 并返回中断标记。
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 final boolean acquireQueued (final Node node, int arg) { boolean failed = true ; try { boolean interrupted = false ; for (;;) { final Node p = node.predecessor(); if (p == head && tryAcquire(arg)) { setHead(node); p.next = null ; failed = false ; return interrupted; } if (shouldParkAfterFailedAcquire(p, node) && parkAndCheckInterrupt()) interrupted = true ; } } finally { if (failed) cancelAcquire(node); } } private static boolean shouldParkAfterFailedAcquire (Node pred, Node node) { int ws = pred.waitStatus; if (ws == Node.SIGNAL) return true ; if (ws > 0 ) { do { node.prev = pred = pred.prev; } while (pred.waitStatus > 0 ); pred.next = node; } else { compareAndSetWaitStatus(pred, ws, Node.SIGNAL); } return false ; }
容易忽略的细节:acquire 不响应中断 。parkAndCheckInterrupt() 只把中断标记记下来继续排队,等真正拿到资源后才 selfInterrupt() 补上;想响应中断要用 acquireInterruptibly,它检测到中断直接抛 InterruptedException 并 cancelAcquire。
2.4 独占模式与共享模式的差异
对比项
独占模式 EXCLUSIVE
共享模式 SHARED
获取方法
acquire / tryAcquire 返回 boolean
acquireShared / tryAcquireShared 返回 int
成功语义
只有当前线程能持有
多个线程可同时持有
节点标记
nextWaiter == null
nextWaiter == SHARED
唤醒行为
release 只唤醒 head 的一个后继
releaseShared 唤醒后可能继续向后传播(PROPAGATE)
典型实现
ReentrantLock
Semaphore、CountDownLatch、ReadLock
共享模式的关键在 setHeadAndPropagate:新 head 就位后,若还有剩余资源(tryAcquireShared 返回 > 0)或 head 状态为 PROPAGATE,会继续 doReleaseShared 唤醒下一个节点,形成级联唤醒 。CountDownLatch 计数归零时,所有 await 线程正是靠它一次性全部放行。
2.5 各类同步器如何复用 AQS
同步器
state 的含义
模式
关键 tryAcquire 逻辑
ReentrantLock
0 未锁;N 表示重入 N 次
独占
CAS 改 0→1;若已是自己则 state+1
ReentrantReadWriteLock
高 16 位读计数,低 16 位写重入
读共享/写独占
读锁看写锁是否被占;写锁看 state 是否非 0
Semaphore
剩余许可数
共享
自旋 CAS 把 state 减掉 acquire 的许可数,不够则返回负数
CountDownLatch
未完成的任务计数
共享
只有 state == 0 才返回 1,否则 -1
FutureTask
任务状态(NEW/COMPLETING 等)
共享
任务完成才返回 1
2.6 手写一个不可重入独占锁 看懂上面之后,自己实现一个锁只需四十行。tryAcquire 里”已锁定就直接返回 false”这一句,正是”不可重入”的定义。
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60 61 62 63 64 65 66 67 68 69 70 71 72 73 74 75 76 77 78 79 80 81 82 83 84 import java.util.concurrent.TimeUnit;import java.util.concurrent.locks.AbstractQueuedSynchronizer;import java.util.concurrent.locks.Condition;public class MutexLock implements java .io.Serializable { private static final long serialVersionUID = 1L ; private static class Sync extends AbstractQueuedSynchronizer { @Override protected boolean isHeldExclusively () { return getState() == 1 ; } @Override protected boolean tryAcquire (int acquires) { assert acquires == 1 ; if (compareAndSetState(0 , 1 )) { setExclusiveOwnerThread(Thread.currentThread()); return true ; } return false ; } @Override protected boolean tryRelease (int releases) { assert releases == 1 ; if (getState() == 0 ) throw new IllegalMonitorStateException (); setExclusiveOwnerThread(null ); setState(0 ); return true ; } Condition newCondition () { return new ConditionObject (); } } private final Sync sync = new Sync (); public void lock () { sync.acquire(1 ); } public void unlock () { sync.release(1 ); } public boolean tryLock () { return sync.tryAcquire(1 ); } public boolean isLocked () { return sync.isHeldExclusively(); } public void lockInterruptibly () throws InterruptedException { sync.acquireInterruptibly(1 ); } public boolean tryLock (long timeout, TimeUnit unit) throws InterruptedException { return sync.tryAcquireNanos(1 , unit.toNanos(timeout)); } public Condition newCondition () { return sync.newCondition(); } public static void main (String[] args) throws InterruptedException { MutexLock lock = new MutexLock (); int [] counter = {0 }; Thread[] ts = new Thread [10 ]; for (int i = 0 ; i < ts.length; i++) { ts[i] = new Thread (() -> { for (int j = 0 ; j < 1000 ; j++) { lock.lock(); try { counter[0 ]++; } finally { lock.unlock(); } } }); ts[i].start(); } for (Thread t : ts) t.join(); System.out.println("counter = " + counter[0 ]); } }
把它改成可重入 只需加一句判断:若 getExclusiveOwnerThread() == Thread.currentThread(),则 setState(getState() + 1) 返回 true,释放时对应减一、减到 0 才真正置 0——这正是 ReentrantLock 的实现。
三、Lock 体系:从接口到实现类 3.1 Lock 接口的方法语义 1 2 3 4 5 6 7 8 public interface Lock { void lock () ; void lockInterruptibly () throws InterruptedException; boolean tryLock () ; boolean tryLock (long time, TimeUnit unit) throws InterruptedException; void unlock () ; Condition newCondition () ; }
对比 synchronized 缺失的能力:tryLock 可做死锁规避 (拿不到第二把锁就先放开第一把,退避重试);lockInterruptibly 让死锁可被外部解除;newCondition 支持一把锁开多个等待队列。
3.2 ReentrantLock 的可重入实现 可重入指”同一线程可多次获取同一把锁”,避免自己锁自己。ReentrantLock 靠两点实现:state 计数 与持有线程判定 :
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 final boolean nonfairTryAcquire (int acquires) { final Thread current = Thread.currentThread(); int c = getState(); if (c == 0 ) { if (compareAndSetState(0 , acquires)) { setExclusiveOwnerThread(current); return true ; } } else if (current == getExclusiveOwnerThread()) { int nextc = c + acquires; if (nextc < 0 ) throw new Error ("Maximum lock count exceeded" ); setState(nextc); return true ; } return false ; } protected final boolean tryRelease (int releases) { int c = getState() - releases; if (Thread.currentThread() != getExclusiveOwnerThread()) throw new IllegalMonitorStateException (); boolean free = false ; if (c == 0 ) { free = true ; setExclusiveOwnerThread(null ); } setState(c); return free; }
3.3 公平锁与非公平锁 公平锁与非公平锁的差别,全部集中在 tryAcquire 的一行代码上。
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 static final class NonfairSync extends Sync { protected boolean tryAcquire (int acquires) { return nonfairTryAcquire(acquires); } } static final class FairSync extends Sync { protected boolean tryAcquire (int acquires) { final Thread current = Thread.currentThread(); int c = getState(); if (c == 0 ) { if (!hasQueuedPredecessors() && compareAndSetState(0 , acquires)) { setExclusiveOwnerThread(current); return true ; } } else if (current == getExclusiveOwnerThread()) { int nextc = c + acquires; if (nextc < 0 ) throw new Error ("Maximum lock count exceeded" ); setState(nextc); return true ; } return false ; } } public final boolean hasQueuedPredecessors () { Node t = tail, h = head, s; return h != t && ((s = h.next) == null || s.thread != Thread.currentThread()); }
一个反直觉的点:公平锁的 tryLock() 依然非公平 ——它直接调 sync.nonfairTryAcquire(1),哪怕对象是 new ReentrantLock(true)。因为”尝试”语义本就尽力而为,要严格公平请用 tryLock(0, TimeUnit.SECONDS)。
性能取舍很清楚:非公平吞吐远高于公平。因为唤醒一个 park 的线程要经历内核态/用户态切换(数千周期) ,公平锁坚持唤醒队首,这段空窗期 CPU 干等;非公平锁允许此刻刚到达、本就在运行态的线程插队,直接把空窗期填满。代价是队尾线程可能饥饿。
为什么默认非公平?因为保证公平的成本高于收益 :非公平吞吐可高出数倍,而饥饿在真实场景极少发生(线程终会执行完)。只有明确要求”先来后到”且临界区较长时,才值得付这笔税。
3.4 Condition:一个锁,多个等待队列 synchronized 的 wait/notify 把等待队列绑定在对象监视器上,一个锁只有一个队列,只能 notifyAll 全量唤醒。AQS 的 ConditionObject 打破了这一限制:一把 ReentrantLock 可以 newCondition() 出任意多个等待队列 ,实现精准唤醒。
3.4.1 await 与 signal 的底层 ConditionObject 内部是一条单向链表 (firstWaiter/lastWaiter),节点复用 Node,waitStatus 为 CONDITION(-2),靠 nextWaiter 串联。流程如下:
await() :addConditionWaiter() 入条件队列 → fullyRelease(node) 完全释放 锁(可重入时 state 可能 >1,必须一次清 0)→ park 挂起 → 被 signal 后由 transferAfterCancelledWait 判定迁移方式 → acquireQueued 重新排队争锁并恢复原 state。
signal() :把条件队列首节点 transferForSignal 到同步队列(状态 CONDITION→0、CAS 入队尾、前驱改 SIGNAL),该线程随后被正常唤醒。
signalAll() :条件队列所有节点依次搬回同步队列。
transferAfterCancelledWait 负责区分”等待期间被中断”的两种情形:
1 2 3 4 5 6 7 8 9 10 11 final boolean transferAfterCancelledWait (Node node) { if (compareAndSetWaitStatus(node, Node.CONDITION, 0 )) { enq(node); return true ; } while (!isOnSyncQueue(node)) Thread.yield (); return false ; }
这个返回值决定了 await() 抛 InterruptedException 的时机:中断早于 signal 才抛异常 ,中断晚于 signal 则先恢复锁、再补上中断标记,保证不丢事件。这比 Object.wait() 的语义更严谨。
3.4.2 用 Condition 实现有界缓冲区 1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60 61 62 63 64 65 import java.util.concurrent.locks.Condition;import java.util.concurrent.locks.ReentrantLock;public class BoundedBuffer <T> { private final Object[] items; private int putPtr, takePtr, count; private final ReentrantLock lock = new ReentrantLock (); private final Condition notFull = lock.newCondition(); private final Condition notEmpty = lock.newCondition(); @SuppressWarnings("unchecked") public BoundedBuffer (int capacity) { items = new Object [capacity]; } public void put (T x) throws InterruptedException { lock.lock(); try { while (count == items.length) notFull.await(); items[putPtr] = x; if (++putPtr == items.length) putPtr = 0 ; count++; notEmpty.signal(); } finally { lock.unlock(); } } @SuppressWarnings("unchecked") public T take () throws InterruptedException { lock.lock(); try { while (count == 0 ) notEmpty.await(); T x = (T) items[takePtr]; if (++takePtr == items.length) takePtr = 0 ; count--; notFull.signal(); return x; } finally { lock.unlock(); } } public static void main (String[] args) { BoundedBuffer<Integer> buf = new BoundedBuffer <>(5 ); new Thread (() -> { for (int i = 1 ; i <= 10 ; i++) { try { buf.put(i); System.out.println("put " + i); Thread.sleep(300 ); } catch (InterruptedException e) { Thread.currentThread().interrupt(); } } }, "producer" ).start(); new Thread (() -> { for (int i = 1 ; i <= 10 ; i++) { try { System.out.println("take " + buf.take()); Thread.sleep(800 ); } catch (InterruptedException e) { Thread.currentThread().interrupt(); } } }, "consumer" ).start(); } }
对比项
Object.wait/notify
Condition.await/signal
依赖
必须先 synchronized 拿到监视器
必须先 lock() 拿到 Lock
等待队列数
每个对象仅 1 个
每把锁可有多个 Condition
唤醒粒度
notify 随机一个 / notifyAll 全部
signal 指定队列的首个 / signalAll 该队列全部
中断语义
抛异常,区分不了中断与 signal 先后
transferAfterCancelledWait 精确区分
超时等待
支持 wait(timeout)
支持 awaitNanos / awaitUntil(绝对时间)
释放方式
释放一次监视器
fullyRelease 释放全部重入次数
3.5 ReentrantReadWriteLock 与锁降级 读写锁把”读-读”从互斥中解放出来:读读共享、读写互斥、写写互斥 ,适合缓存、配置中心这类读远多于写的场景。
锁降级指”持有写锁 → 获取读锁 → 释放写锁 “,最终降级为读锁,保证刚写完的数据立刻能被自己读到,且期间不会被其它写线程插入。
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 import java.util.concurrent.locks.ReentrantReadWriteLock;public class CacheDemo { private final ReentrantReadWriteLock rwl = new ReentrantReadWriteLock (); private final ReentrantReadWriteLock.ReadLock readLock = rwl.readLock(); private final ReentrantReadWriteLock.WriteLock writeLock = rwl.writeLock(); private volatile boolean cacheValid; private Object data; public Object getData () { readLock.lock(); if (!cacheValid) { readLock.unlock(); writeLock.lock(); try { if (!cacheValid) { data = loadFromDb(); cacheValid = true ; } readLock.lock(); } finally { writeLock.unlock(); } } try { return data; } finally { readLock.unlock(); } } private Object loadFromDb () { System.out.println(Thread.currentThread().getName() + " 加载数据库..." ); return new Object (); } }
为什么不支持锁升级(先读后写) ?读锁可能被多线程持有,若允许其中一个升级成写锁,其它读线程还浑然不觉地持有读锁,互斥语义即被破坏;且两个读线程同时升级会互相等待而死锁。所以要先释放再获取,或一开始就拿写锁。
3.6 StampedLock:乐观读 StampedLock(JDK 8)用 long stamp 作票据,提供三种模式:
模式
方法
语义
是否阻塞
写锁
writeLock() / unlockWrite(stamp)
独占
是
悲观读锁
readLock() / unlockRead(stamp)
共享
是
乐观读
tryOptimisticRead() / validate(stamp)
不加锁,读完校验是否被写过
否
乐观读的思路类似数据库的乐观锁版本号:读前取 stamp,读完用 validate(stamp) 检查期间有无写操作,没有就直接用,有则升级为悲观读锁重读一次。
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 import java.util.concurrent.locks.StampedLock;public class Point { private double x, y; private final StampedLock sl = new StampedLock (); public void move (double deltaX, double deltaY) { long stamp = sl.writeLock(); try { x += deltaX; y += deltaY; } finally { sl.unlockWrite(stamp); } } public double distanceFromOrigin () { long stamp = sl.tryOptimisticRead(); double currentX = x, currentY = y; if (!sl.validate(stamp)) { stamp = sl.readLock(); try { currentX = x; currentY = y; } finally { sl.unlockRead(stamp); } } return Math.sqrt(currentX * currentX + currentY * currentY); } public double pessimisticDistance () { long stamp = sl.readLock(); try { return Math.sqrt(x * x + y * y); } finally { sl.unlockRead(stamp); } } }
StampedLock 有两个坑:不可重入 (同线程重复获取会自锁)、不支持 Condition (需要条件等待请退回 ReentrantReadWriteLock);且它的悲观读并非 AQS 实现,不能当普通 Lock 传给需要 Lock 接口的方法。
3.7 死锁排查:jstack ReentrantLock 死锁不会像 synchronized 那样被 jstack 自动标注 “Found one Java-level deadlock”,但会打印 WAITING (parking) 与 “parking to wait for <0x...> (a ReentrantLock$NonfairSync)”,需人工串联。
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 import java.util.concurrent.locks.ReentrantLock;import java.util.concurrent.TimeUnit;public class LockDeadlockDemo { private static final ReentrantLock lockA = new ReentrantLock (); private static final ReentrantLock lockB = new ReentrantLock (); public static void main (String[] args) throws InterruptedException { new Thread (() -> { lockA.lock(); try { TimeUnit.SECONDS.sleep(1 ); System.out.println("T1 尝试获取 lockB..." ); lockB.lock(); try { System.out.println("T1 拿到两把锁" ); } finally { lockB.unlock(); } } catch (InterruptedException e) { Thread.currentThread().interrupt(); } finally { lockA.unlock(); } }, "Thread-A-B" ).start(); new Thread (() -> { lockB.lock(); try { TimeUnit.SECONDS.sleep(1 ); System.out.println("T2 尝试获取 lockA..." ); lockA.lock(); try { System.out.println("T2 拿到两把锁" ); } finally { lockA.unlock(); } } catch (InterruptedException e) { Thread.currentThread().interrupt(); } finally { lockB.unlock(); } }, "Thread-B-A" ).start(); TimeUnit.SECONDS.sleep(3 ); System.out.println("主线程退出,请用 jstack <pid> 观察两个 WAITING (parking) 线程" ); } }
排查步骤:jps -l 取进程号 → jstack -l <pid> > dead.txt → 搜 WAITING (parking) 记下 parking to wait for <0x...> → 再搜该地址出现在哪个线程的 locked <0x...> 之后,即可画出”谁持有了谁想要的锁”的环。修复手段通常是固定加锁顺序 (按对象 hashCode 排序加锁),或用 tryLock(timeout) + 退避重试。
四、原子类家族:无锁化的第一选择 4.1 CAS 原理 CAS(Compare And Swap)是一条 CPU 原子指令:给定内存位置 V、期望旧值 A、新值 B,当且仅当 V 等于 A 时才把 V 改成 B,否则什么都不做 ,并返回是否成功。x86 上对应 cmpxchg,配合 lock 前缀锁总线/缓存行保证原子性。Java 层通过 sun.misc.Unsafe 暴露:
1 2 3 4 public final native boolean compareAndSwapObject (Object o, long offset, Object expected, Object x) ;public final native boolean compareAndSwapInt (Object o, long offset, int expected, int x) ;public final native boolean compareAndSwapLong (Object o, long offset, long expected, long x) ;
JDK 9 之后 Unsafe 被逐步收口,官方替代品是 VarHandle(JDK 9)与 MemorySegment(JDK 22 的 FFM API),AtomicInteger 内部也换成了 VarHandle,语义完全一致。
4.2 CAS 的三大问题与解法
问题
描述
解决方案
ABA
值从 A 变成 B 又变回 A,CAS 检查时以为没变过
AtomicStampedReference(版本号)、AtomicMarkableReference(布尔标记)
循环开销大
竞争激烈时 CAS 长期失败,自旋空耗 CPU
LongAdder 分散热点;或退化为锁;JVM 支持 pause 指令降低功耗
只能保证单变量原子性
多个变量需要一起原子更新时无能为力
封装成对象用 AtomicReference;或直接用锁
ABA 的经典场景:栈顶元素从 A 弹出换成 B 又被压回 A,某线程 CAS 判断”栈顶还是 A”而成功,但整个栈的结构其实已经变了。
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 import java.util.concurrent.atomic.AtomicStampedReference;public class ABADemo { public static void main (String[] args) throws InterruptedException { AtomicStampedReference<Integer> ref = new AtomicStampedReference <>(100 , 1 ); Thread t1 = new Thread (() -> { int stamp = ref.getStamp(); System.out.println("t1 读到 stamp=" + stamp); try { Thread.sleep(1000 ); } catch (InterruptedException e) { return ; } boolean ok = ref.compareAndSet(100 , 101 , stamp, stamp + 1 ); System.out.println("t1 CAS 结果 = " + ok + ",当前值 = " + ref.getReference() + ",当前版本 = " + ref.getStamp()); }, "t1" ); Thread t2 = new Thread (() -> { ref.compareAndSet(100 , 101 , ref.getStamp(), ref.getStamp() + 1 ); System.out.println("t2 第一次修改,值=" + ref.getReference() + " 版本=" + ref.getStamp()); ref.compareAndSet(101 , 100 , ref.getStamp(), ref.getStamp() + 1 ); System.out.println("t2 第二次修改(改回 100),值=" + ref.getReference() + " 版本=" + ref.getStamp()); }, "t2" ); t1.start(); Thread.sleep(100 ); t2.start(); t1.join(); t2.join(); } }
4.3 原子类家族全景
分类
类
用途
基本类型
AtomicInteger、AtomicLong、AtomicBoolean
单变量原子读写与运算
数组
AtomicIntegerArray、AtomicLongArray、AtomicReferenceArray
数组元素的原子更新(复制了数组,不影响原数组)
引用类型
AtomicReference、AtomicStampedReference、AtomicMarkableReference
对象引用的原子更新,后两者解决 ABA
字段更新器
AtomicIntegerFieldUpdater、AtomicLongFieldUpdater、AtomicReferenceFieldUpdater
反射式更新对象的 volatile 字段,省去包装对象
累加器(JDK 8)
LongAdder、DoubleAdder、LongAccumulator、DoubleAccumulator
高并发求和/自定义聚合,吞吐远高于原子类
其它
AtomicLongFieldUpdater 等
见上;Striped64 是 LongAdder 的父类
4.4 AtomicInteger 的 incrementAndGet 1 2 3 4 5 6 7 8 9 10 11 12 13 public final int incrementAndGet () { return unsafe.getAndAddInt(this , valueOffset, 1 ) + 1 ; } public final int getAndAddInt (Object o, long offset, int delta) { int v; do { v = this .getIntVolatile(o, offset); } while (!this .compareAndSwapInt(o, offset, v, v + delta)); return v; }
JDK 9+ 语义相同,只是把 Unsafe 换成 VarHandle.compareAndSet,并在失败循环里加 Thread.onSpinWait() 提示 CPU(x86 上是 pause 指令),降低自旋功耗。
4.5 LongAdder:分段累加与伪共享 AtomicLong 高并发累加时有致命弱点:所有线程 CAS 同一个变量 ,只有一个能成功,其余全部自旋重试,竞争度随线程数线性恶化。
LongAdder 的思路是”分而治之”,继承自 Striped64:
base :竞争不激烈时直接 CAS 这个基础值,和 AtomicLong 一样快;
Cell[] :竞争激烈时每线程按自己的 probe 哈希映射到某个 Cell 槽位各加各的,冲突就 rehash 换槽或扩容数组(上限为 CPU 核数);
sum() :base + 所有 Cell 求和。这不是原子快照 ,并发更新时可能漏掉增量,所以 LongAdder 只适合 QPS、计数这类统计场景,不适合做需要精确一致性的状态判断。
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 @sun .misc.Contended static final class Cell { volatile long value; Cell(long x) { value = x; } final boolean cas (long cmp, long val) { return UNSAFE.compareAndSwapLong(this , valueOffset, cmp, val); } } public void add (long x) { Cell[] cs; long b, v; int m; Cell c; if ((cs = cells) != null || !casBase(b = base, b + x)) { boolean uncontended = true ; if (cs == null || (m = cs.length - 1 ) < 0 || (c = cs[getProbe() & m]) == null || !(uncontended = c.cas(v = c.value, v + x))) longAccumulate(x, null , uncontended); } }
@Contended 是点睛之笔。CPU 以 Cache Line(通常 64 字节) 加载内存,一个 Cell 仅 16 字节,相邻两个 Cell 会落在同一条缓存行:A 改 Cell[0]、B 改 Cell[1] 虽逻辑无关,MESI 协议却会互相把对方的缓存行置为无效,产生伪共享(False Sharing) ,性能退化到与单变量 CAS 相当。@Contended 让 JVM 插入 128 字节填充(默认只对 JDK 内部类生效,用户类需 -XX:-RestrictContended),使每个 Cell 独占一条缓存行。
实测对比(4 核机器,32 线程各累加 1000 万次,量级仅供参考):
实现方式
耗时(约)
相对吞吐
结果准确性
synchronized 方法
4.5 s
1x
精确
ReentrantLock
1.8 s
2.5x
精确
AtomicLong
1.2 s
3.7x
精确
LongAdder
0.25 s
18x
最终一致(sum() 非原子快照)
结论很直接:高并发计数用 LongAdder,需要精确值或做 CAS 判断用 AtomicLong,低并发下 AtomicLong 更快也更省内存 (后者要维护 Cell 数组)。LongAccumulator 是通用版,可传入自定义二元运算(如 Long::max)与初始值。
五、并发容器:性能与语义的权衡 5.1 同步容器 vs 并发容器 同步容器(Vector、Hashtable、Collections.synchronizedXxx)在每个方法上加 synchronized 锁住整个容器,并发容器则把锁粒度降到元素级:
对比项
同步容器(Vector、Hashtable、Collections.synchronizedXxx)
并发容器(ConcurrentHashMap、CopyOnWriteArrayList 等)
实现方式
方法上加 synchronized,锁住整个容器
分段/CAS/写时复制,锁粒度极细
并发度
同一时刻仅一个线程可访问
读读、读写(多数情况)可并发
复合操作
size() 与 get() 之间需外部加锁
ConcurrentHashMap 提供 putIfAbsent 等原子复合方法
迭代器
快速失败(fail-fast),并发修改抛 ConcurrentModificationException
弱一致(fail-safe),迭代期间允许修改
性能
低,高竞争下急剧退化
高,随线程数近似线性扩展
适用场景
遗留代码、极低并发
一切新代码
5.2 ConcurrentHashMap:JDK 8 的重写 JDK 7 用 Segment 分段锁 (继承 ReentrantLock,默认 16 段),并发度上限就是段数。JDK 8 彻底重写,改为 数组 + 链表/红黑树 + CAS + synchronized 锁单个桶 :
put :桶空时 CAS 插入(无锁);桶非空则 synchronized 锁住桶的头节点,只锁一条链表,其它桶不受影响;
get :全程无锁,Node 的 val、next 均为 volatile,靠 volatile 读保证可见性;
计数 :借鉴 LongAdder,用 baseCount + CounterCell[] 分段计数,size() 为求和结果(非精确),推荐 mappingCount()(long);
扩容 :transfer 按 stride 把迁移任务切块,线程做完自己那段若还有未迁移区间会继续领取,多线程可协助迁移 (helpTransfer),靠 ForwardingNode(hash = MOVED = -1)标记该桶已搬走,扩容期间读写仍可正常进行;
树化 :链表长度 ≥ 8 且数组长度 ≥ 64 时转红黑树,≤ 6 时退化回链表。
为什么不允许 null 键值 ?为了消除歧义:get(key) 返回 null 时无法区分”key 不存在”与”key 存在但值为 null”。单线程下还能用 containsKey 二次确认,并发下两次调用之间状态已变,判断永远不可靠。所以 Doug Lea 直接在设计上禁止 null,让 null 唯一表示”不存在”。
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 import java.util.concurrent.ConcurrentHashMap;import java.util.concurrent.atomic.LongAdder;public class ConcurrentHashMapDemo { public static void main (String[] args) { ConcurrentHashMap<String, LongAdder> wordCount = new ConcurrentHashMap <>(); wordCount.putIfAbsent("java" , new LongAdder ()); wordCount.get("java" ).increment(); wordCount.computeIfAbsent("juc" , k -> new LongAdder ()).increment(); wordCount.computeIfAbsent("juc" , k -> new LongAdder ()).increment(); wordCount.forEach((k, v) -> System.out.println(k + " -> " + v.sum())); System.out.println("总 key 数 = " + wordCount.mappingCount()); System.out.println("总次数 = " + wordCount.reduceValuesToLong(1L , LongAdder::sum)); } }
5.3 CopyOnWriteArrayList 写时复制:每次修改(add/set/remove)都复制一份新数组 ,改完后用 volatile 引用切换过去,读操作完全不加锁。
代价显而易见:写期间新旧两个数组同时在堆里,元素多时易触发 GC;复制是 O(n),写多即灾难;迭代器持有创建时的快照 ,遍历期间其它线程的修改完全不可见,即”弱一致性”。
适用场景很窄:读极多、写极少且能容忍短暂不一致 ,典型是监听器列表、路由表、黑白名单配置。
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 import java.util.concurrent.CopyOnWriteArrayList;public class CopyOnWriteDemo { public static void main (String[] args) throws InterruptedException { CopyOnWriteArrayList<String> listeners = new CopyOnWriteArrayList <>(); listeners.add("listener-1" ); listeners.add("listener-2" ); Thread writer = new Thread (() -> { try { Thread.sleep(10 ); } catch (InterruptedException e) { return ; } listeners.add("listener-3" ); System.out.println("写线程已添加 listener-3" ); }); writer.start(); for (String s : listeners) { System.out.println("遍历到:" + s); Thread.sleep(50 ); } writer.join(); System.out.println("遍历结束后 size = " + listeners.size()); } }
5.4 ConcurrentLinkedQueue 基于 Michael-Scott 算法的无界非阻塞队列 :入队 CAS 尾节点的 next,出队 CAS head 并帮助推进,全程无锁,适合高并发且不需要阻塞语义的场景。注意 size() 需遍历,是 O(n) 操作 ,判空请用 isEmpty()。
5.5 BlockingQueue 家族与选型
队列
底层结构
是否有界
锁/实现
典型用途
ArrayBlockingQueue
数组
有界 (构造指定)
单把 ReentrantLock + 两个 Condition
固定容量、需要背压的池化场景
LinkedBlockingQueue
链表
可选(默认 Integer.MAX_VALUE)
双锁(putLock/takeLock)分离,吞吐更高
通用任务队列,线程池默认项
PriorityBlockingQueue
堆
无界(会自动扩容)
单锁 + 自旋 CAS 扩容
需要按优先级出队的任务调度
DelayQueue
PriorityQueue
无界
单锁 + Condition available
延时任务、订单超时、缓存过期
SynchronousQueue
无存储
容量为 0
CAS 双栈/双队列
直接交接,Executors.newCachedThreadPool 用
LinkedTransferQueue
链表
无界
CAS + transfer 语义
生产者需确认”已被消费者接收”
LinkedBlockingDeque
双向链表
可选
单锁
工作窃取、双端队列
DelayedWorkQueue
堆(数组)
有界/自扩容
ScheduledThreadPoolExecutor 内部专用
定时任务调度
四组插入/移除 API 的语义必须分清:
行为
抛异常
返回特殊值
阻塞
超时
插入
add(e)
offer(e)
put(e)
offer(e, time, unit)
移除
remove()
poll()
take()
poll(time, unit)
检查
element()
peek()
—
—
5.5.1 用 DelayQueue 实现订单超时关闭 1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 import java.util.concurrent.*;public class OrderTimeoutDemo { static class OrderTask implements Delayed { private final String orderId; private final long expireAt; public OrderTask (String orderId, long delayMillis) { this .orderId = orderId; this .expireAt = System.currentTimeMillis() + delayMillis; } @Override public long getDelay (TimeUnit unit) { return unit.convert(expireAt - System.currentTimeMillis(), TimeUnit.MILLISECONDS); } @Override public int compareTo (Delayed o) { return Long.compare(this .expireAt, ((OrderTask) o).expireAt); } public String getOrderId () { return orderId; } } public static void main (String[] args) throws InterruptedException { DelayQueue<OrderTask> delayQueue = new DelayQueue <>(); delayQueue.put(new OrderTask ("ORDER-001" , 2_000 )); delayQueue.put(new OrderTask ("ORDER-002" , 5_000 )); delayQueue.put(new OrderTask ("ORDER-003" , 3_000 )); System.out.println("三个订单已下单,等待超时关闭..." ); Thread closer = new Thread (() -> { while (!Thread.currentThread().isInterrupted()) { try { OrderTask task = delayQueue.take(); System.out.println("[" + System.currentTimeMillis() % 100_000 + "] 订单 " + task.getOrderId() + " 超时未支付,执行关闭" ); } catch (InterruptedException e) { Thread.currentThread().interrupt(); break ; } } }, "order-timeout-closer" ); closer.setDaemon(true ); closer.start(); closer.join(6_000 ); } }
生产环境补充:DelayQueue 是单机内存 方案,重启即丢失。更可靠的做法是延迟消息(RocketMQ/RabbitMQ 死信或 Redis ZSet 轮询)+ 数据库兜底扫描,或接入分布式调度(XXL-Job、ElasticJob)。
5.5.2 生产者消费者的三种写法 1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60 61 62 63 64 65 66 67 68 69 70 71 72 73 74 75 76 77 78 79 80 81 82 83 84 85 86 87 88 89 90 91 92 93 94 95 96 import java.util.concurrent.*;import java.util.concurrent.atomic.AtomicInteger;public class ProducerConsumerDemo { static void style1 () { BlockingQueue<String> q = new ArrayBlockingQueue <>(10 ); new Thread (() -> { try { for (int i = 0 ; i < 5 ; i++) { q.put("item-" + i); System.out.println("生产 item-" + i); } } catch (InterruptedException e) { Thread.currentThread().interrupt(); } }).start(); new Thread (() -> { try { for (int i = 0 ; i < 5 ; i++) { System.out.println("消费 " + q.take()); } } catch (InterruptedException e) { Thread.currentThread().interrupt(); } }).start(); } static void style2 () { BlockingQueue<String> q = new LinkedBlockingQueue <>(); volatile boolean [] running = {true }; new Thread (() -> { while (running[0 ]) { try { String item = q.poll(1 , TimeUnit.SECONDS); if (item != null ) System.out.println("消费 " + item); } catch (InterruptedException e) { Thread.currentThread().interrupt(); break ; } } System.out.println("消费者退出" ); }).start(); try { q.offer("a" , 2 , TimeUnit.SECONDS); q.offer("b" , 2 , TimeUnit.SECONDS); Thread.sleep(1200 ); running[0 ] = false ; } catch (InterruptedException e) { Thread.currentThread().interrupt(); } } static void style3 () throws InterruptedException { int producers = 2 , consumers = 3 ; BlockingQueue<String> q = new LinkedBlockingQueue <>(20 ); String POISON = "POISON_PILL" ; AtomicInteger remainingConsumers = new AtomicInteger (consumers); CountDownLatch done = new CountDownLatch (consumers); for (int c = 0 ; c < consumers; c++) { new Thread (() -> { try { while (true ) { String item = q.take(); if (POISON.equals(item)) { if (remainingConsumers.decrementAndGet() > 0 ) q.put(POISON); break ; } System.out.println(Thread.currentThread().getName() + " 消费 " + item); } } catch (InterruptedException e) { Thread.currentThread().interrupt(); } finally { done.countDown(); } }, "consumer-" + c).start(); } for (int p = 0 ; p < producers; p++) { final int pid = p; new Thread (() -> { for (int i = 0 ; i < 6 ; i++) { try { q.put("P" + pid + "-item" + i); } catch (InterruptedException e) { return ; } } }, "producer-" + p).start(); } Thread.sleep(500 ); q.put(POISON); done.await(); System.out.println("所有消费者已优雅退出" ); } public static void main (String[] args) throws InterruptedException { style1(); Thread.sleep(500 ); System.out.println("--- 写法二 ---" ); style2(); Thread.sleep(1500 ); System.out.println("--- 写法三 ---" ); style3(); } }
六、同步协作工具类 6.1 四工具对照
工具
是否可重用
计数方向
是否阻塞等待
核心语义
典型场景
CountDownLatch
否 (一次性)
递减到 0
await() 阻塞
一个或多个线程等其它线程做完
启动检查、并发压测、服务就绪
CyclicBarrier
是 (可 reset)
递增到 parties
await() 阻塞
一组线程互相等待,到齐一起放行
多阶段计算、并行迭代
Semaphore
是
许可加减
acquire() 阻塞
控制同时访问的线程数
限流、资源池、数据库连接数
Phaser
是(动态注册)
分阶段(phase)
arriveAndAwaitAdvance
多阶段 + 参与者可动态增减
复杂流水线、替代 Barrier+Latch
Exchanger
是
成对交换
exchange() 阻塞
两个线程在汇合点交换数据
双缓冲、校对数据
6.2 CountDownLatch:一次性倒数门闩 内部是 AQS 的共享模式 :state 即计数值,await() 等价于 acquireSharedInterruptibly(1),只有 state 归零才返回 1;countDown() 等价于 releaseShared(1),减到 0 时触发级联唤醒。归零后无法重置 ,是一次性消耗品。
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60 import java.util.concurrent.*;import java.util.concurrent.atomic.AtomicInteger;public class ConcurrentPressureTest { public static void main (String[] args) throws InterruptedException { int threadCount = 100 ; ExecutorService pool = Executors.newFixedThreadPool(32 ); CountDownLatch startGate = new CountDownLatch (1 ); CountDownLatch endGate = new CountDownLatch (threadCount); AtomicInteger success = new AtomicInteger (); AtomicInteger failure = new AtomicInteger (); AtomicLongHolder totalCost = new AtomicLongHolder (); for (int i = 0 ; i < threadCount; i++) { pool.submit(() -> { try { startGate.await(); long start = System.currentTimeMillis(); boolean ok = mockCallApi(); long cost = System.currentTimeMillis() - start; totalCost.add(cost); if (ok) success.incrementAndGet(); else failure.incrementAndGet(); } catch (InterruptedException e) { Thread.currentThread().interrupt(); } finally { endGate.countDown(); } }); } long begin = System.currentTimeMillis(); startGate.countDown(); endGate.await(); long wall = System.currentTimeMillis() - begin; System.out.println("线程数 = " + threadCount); System.out.println("成功 = " + success.get() + ",失败 = " + failure.get()); System.out.println("总耗时 = " + wall + " ms" ); System.out.println("平均 RT = " + (totalCost.sum() / threadCount) + " ms" ); System.out.println("估算 QPS = " + (threadCount * 1000L / wall)); pool.shutdown(); } private static boolean mockCallApi () throws InterruptedException { TimeUnit.MILLISECONDS.sleep(50 + (long ) (Math.random() * 50 )); return Math.random() > 0.05 ; } static class AtomicLongHolder { private final LongAdder adder = new LongAdder (); void add (long v) { adder.add(v); } long sum () { return adder.sum(); } } }
6.3 CyclicBarrier:可复用的栅栏 与 CountDownLatch 的两点本质区别:可循环使用 (一代结束自动重置,也可 reset()),以及支持 barrierAction (所有线程到齐后、放行前,由最后一个到达的线程执行的 Runnable)。内部靠 ReentrantLock + Condition 实现,而非 AQS 共享模式。
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 import java.util.concurrent.*;public class CyclicBarrierDemo { public static void main (String[] args) { int parties = 3 ; int [][] matrix = { {1 , 2 , 3 }, {4 , 5 , 6 }, {7 , 8 , 9 } }; int [] partial = new int [parties]; CyclicBarrier barrier = new CyclicBarrier (parties, () -> { int sum = 0 ; for (int p : partial) sum += p; System.out.println("本阶段汇总结果 = " + sum); }); ExecutorService pool = Executors.newFixedThreadPool(parties); for (int i = 0 ; i < parties; i++) { final int idx = i; pool.submit(() -> { try { for (int phase = 0 ; phase < 2 ; phase++) { int s = 0 ; for (int v : matrix[idx]) s += v + phase; partial[idx] = s; System.out.println(Thread.currentThread().getName() + " 第 " + phase + " 阶段完成,部分和=" + s); barrier.await(); } } catch (InterruptedException | BrokenBarrierException e) { System.err.println("栅栏被打破:" + e.getMessage()); } }); } pool.shutdown(); } }
6.4 Semaphore:许可模型与限流器 Semaphore 同样是 AQS 共享模式,state 即剩余许可数。acquire(n) 拿 n 个许可,release(n) 归还(可多于 acquire 的数量 ,等于动态扩容许可)。典型用途是限流与对象池。
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60 61 62 63 64 65 66 67 68 69 70 71 import java.util.concurrent.*;import java.util.concurrent.atomic.AtomicInteger;public class SemaphoreLimiter { private final Semaphore semaphore; private final boolean fair; public SemaphoreLimiter (int permits , boolean fair) { this .semaphore = new Semaphore (permits , fair); this .fair = fair; } public <T> T execute (Callable<T> task) throws Exception { semaphore.acquire(); try { return task.call(); } finally { semaphore.release(); } } public <T> T execute (Callable<T> task, long timeout, TimeUnit unit) throws Exception { if (!semaphore.tryAcquire(timeout, unit)) { throw new RejectedExecutionException ("系统繁忙,请稍后再试" ); } try { return task.call(); } finally { semaphore.release(); } } public int availablePermits () { return semaphore.availablePermits(); } public int queueLength () { return semaphore.getQueueLength(); } public static void main (String[] args) throws InterruptedException { SemaphoreLimiter limiter = new SemaphoreLimiter (5 , true ); AtomicInteger ok = new AtomicInteger (), rejected = new AtomicInteger (); CountDownLatch latch = new CountDownLatch (50 ); ExecutorService pool = Executors.newFixedThreadPool(50 ); for (int i = 0 ; i < 50 ; i++) { final int req = i; pool.submit(() -> { try { limiter.execute(() -> { TimeUnit.MILLISECONDS.sleep(200 ); return "req-" + req + " done" ; }, 1 , TimeUnit.SECONDS); ok.incrementAndGet(); } catch (RejectedExecutionException e) { rejected.incrementAndGet(); } catch (Exception e) { e.printStackTrace(); } finally { latch.countDown(); } }); } latch.await(); System.out.println("成功 = " + ok.get() + ",被限流 = " + rejected.get()); System.out.println("剩余许可 = " + limiter.availablePermits() + ",排队线程 = " + limiter.queueLength()); pool.shutdown(); } }
6.5 Phaser 与 Exchanger 简述 Phaser 可看作 CountDownLatch + CyclicBarrier 的增强版,支持动态注册/注销参与者 (register/arriveAndDeregister)与多阶段 (getPhase()),重写 onAdvance 可做阶段回调并返回 true 终止,适合分阶段流水线。
Exchanger 是双线程汇合点:exchange(V x) 阻塞直到另一线程也调用它,然后两者交换数据并返回对方的值,适用于双缓冲、数据校对。注意它两两配对 ,奇数个线程会有一个永远等待(除非设超时)。
七、ThreadLocal 深入 7.1 用途与基本用法 ThreadLocal 提供线程封闭 :每个线程持有变量的独立副本,常用于隐式传参(用户上下文、TraceId、事务上下文)。代价是引入了一条隐藏的调用链依赖,排查时不易追踪。
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 import java.util.concurrent.ExecutorService;import java.util.concurrent.Executors;public class ThreadLocalBasic { private static final ThreadLocal<UserContext> CONTEXT = new ThreadLocal <>(); static class UserContext { final String userId; UserContext(String userId) { this .userId = userId; } } public static void main (String[] args) { ExecutorService pool = Executors.newFixedThreadPool(2 ); for (int i = 1 ; i <= 4 ; i++) { final String uid = "user-" + i; pool.submit(() -> { try { CONTEXT.set(new UserContext (uid)); service(); } finally { CONTEXT.remove(); } }); } pool.shutdown(); } static void service () { dao(); } static void dao () { UserContext ctx = CONTEXT.get(); System.out.println(Thread.currentThread().getName() + " 执行 SQL,uid=" + ctx.userId); } }
7.2 ThreadLocalMap 的结构 ThreadLocal 本身不存值,真正的容器是每个 Thread 内部的 ThreadLocalMap (Thread.threadLocals 字段),key 是 ThreadLocal 实例自身,value 是业务值。
1 2 3 4 5 6 7 8 9 10 11 static class ThreadLocalMap { static class Entry extends WeakReference <ThreadLocal<?>> { Object value; Entry(ThreadLocal<?> k, Object v) { super (k); value = v; } } private Entry[] table; }
7.3 内存泄漏链条与”为什么 value 不是弱引用” Entry 的 key 是弱引用 :外部对 ThreadLocal 实例的强引用消失后,下次 GC 就回收 key,Entry 变成 key == null 的脏条目。但 value 仍是强引用 ,被 Thread → threadLocals → Entry → value 牢牢拽着。若线程不死亡(线程池核心线程几乎永不死亡 ),value 就永不回收——这就是完整的泄漏链条:
1 Thread(长生命周期) → ThreadLocalMap → Entry(key=null) → value(强引用,泄漏)
虽然 ThreadLocalMap 在 set/get/remove 时会顺带做启发式清理(expungeStaleEntry 清掉 key 为 null 的槽位),但这是被动的:之后若不再访问这个 Map,脏条目就一直躺着。
那为什么不把 value 也设成弱引用? 因为那样更糟:value 的唯一强引用就来自 Entry,设为弱引用一次 GC 就可能把你还在用的值清成 null,产生比泄漏更难排查的 bug。key 用弱引用是因为其生命周期由外部持有者决定(通常 static final);value 只能靠 remove() 主动断开。
7.4 线程池中的脏数据问题 比泄漏更常见的是脏数据 :线程池复用线程,上次请求设的值没清理,下次请求 get() 直接拿到上一个用户的数据——严重时会造成跨用户数据泄露。三条铁律:try-finally 包裹、remove() 放 finally;入口统一设置、出口统一清理 (Filter/Interceptor/AOP 最合适);不在异步子线程里直接读父线程的 ThreadLocal。
7.5 InheritableThreadLocal 与 TransmittableThreadLocal
类型
父子线程传递
线程池场景
原理
ThreadLocal
不传递
不可用
数据存在各自 Thread 的 map 里
InheritableThreadLocal
创建子线程时 拷贝一次
不可用 (线程池线程早已创建)
Thread 构造时把父 inheritableThreadLocals 复制到子线程
TransmittableThreadLocal(TTL)
提交任务时传递
可用
阿里开源,包装 Runnable,在任务提交时刻 捕获上下文,执行前注入、执行后还原
InheritableThreadLocal 只在 new Thread() 那一刻生效,而线程池的 worker 线程早已创建并被复用,父线程上下文根本传不进去。TTL 用 TtlRunnable.get(runnable) 包装任务,在 submit 时快照上下文、run 前注入、run 后还原,是目前 TraceId 异步透传的事实标准。
7.6 框架中的典型应用
Spring 事务 :TransactionSynchronizationManager 用多个 ThreadLocal 保存当前事务的 ConnectionHolder,保证同一线程内多个 DAO 拿到同一个连接——这正是事务成立的前提;也因此事务上下文不能跨线程传播 ,子线程的数据库操作不参与主线程事务。
MDC 日志 :MDC 底层是 ThreadLocal<Map<String,String>>,Filter 里 MDC.put("traceId", id)、模板加 %X{traceId} 即可全链路打印,记得 finally 里 MDC.clear()。
八、实战与面试题 8.1 AQS 与锁 Q1:AQS 的核心思想? 三件套:volatile int state 表示资源、CLH 变体的 FIFO 双向队列管排队、Node.waitStatus 表示节点状态;再用模板方法把排队/阻塞/唤醒固化在基类,只把 tryAcquire/tryRelease/tryAcquireShared/tryReleaseShared/isHeldExclusively 留给子类。
Q2:AQS 为什么用 CLH 队列变体? 一是入队必须无锁,CAS 挂尾 + 自旋重试天然适应;二是要支持节点取消,双向 prev 让 cancelAcquire 能摘除中间节点,单向做不到;三是把自旋等待改为 LockSupport.park(),避免空转烧 CPU。
Q3:SIGNAL 有什么用?为什么必须先设前驱为 SIGNAL 才能 park? SIGNAL 表示”我释放时有义务唤醒后继”。不设就 park,前驱释放时不知道后面有人等就不会 unpark,该线程永久挂起。shouldParkAfterFailedAcquire 先 CAS 把前驱改成 SIGNAL 并返回 false 让外层再转一圈,正是为闭合这个契约。
Q4:公平锁和非公平锁的实现差异? 差异只在 tryAcquire:非公平直接 CAS 抢;公平在 CAS 前先调 hasQueuedPredecessors() 判断队中是否有更早的等待者。另注意 ReentrantLock.tryLock() 始终非公平,即便对象创建为公平锁。
Q5:为什么默认非公平锁? 唤醒 park 的线程要内核态切换,期间 CPU 空转;非公平允许此刻刚到达、本就在运行态的线程插队,填满空窗期,吞吐可高出数倍。公平锁坚持先来后到,代价是大量无效唤醒,而饥饿在真实场景极少发生。
Q6:Condition.await() 为什么必须”完全释放”锁? 锁可重入,state 可能大于 1。fullyRelease 一次把 state 清 0 并置空持有线程,否则其它线程永远拿不到锁;唤醒后再用 acquireQueued 抢锁并把 state 恢复到原重入深度。
Q7:transferAfterCancelledWait 的作用? 区分”中断发生在 signal 之前还是之后”。把 CONDITION CAS 成 0 成功,说明还没被 signal(自己超时/中断),需自己 enq 回同步队列,返回 true,await 抛 InterruptedException;失败说明已被 signal 迁移,自旋等其入队即可,返回 false,await 正常返回但补中断标记——既不丢 signal 也不吞中断。
Q8:读写锁为什么支持降级不支持升级? 降级是”写锁 → 读锁 → 释放写锁”,始终自己持有,安全。升级是”读锁 → 写锁”,而读锁可能被多线程持有,允许其中一个升级则互斥语义被破坏,且两个读线程同时升级会互相等待死锁。
Q9:synchronized 和 ReentrantLock 怎么选? 优先 synchronized(不会忘解锁、JVM 有锁消除与自适应自旋);需要可中断、超时、尝试加锁、公平性、多条件队列、读写分离时才上 ReentrantLock。
8.2 原子类与容器 Q10:CAS 的三大问题及解决方案? ABA:用 AtomicStampedReference(版本戳)或 AtomicMarkableReference(布尔标记)。循环开销大:高并发改用 LongAdder 分散热点或退化为锁,JDK 9+ 用 Thread.onSpinWait() 降低自旋功耗。只能保证单变量原子性:把多字段封装成对象用 AtomicReference,或直接用锁。
Q11:LongAdder 为什么比 AtomicLong 快?sum() 精确吗? base + Cell[] 分段累加,各线程哈希到不同 Cell 独立 CAS,把单点竞争拆成多槽竞争,冲突就 rehash 或扩容;Cell 上加 @Contended 填充缓存行避免伪共享。缺点:sum() 是 base 与所有 Cell 求和,并发更新时不是原子快照 ,只适合统计场景。低并发下 AtomicLong 更快更省内存。
Q12:ConcurrentHashMap 在 JDK 7 和 JDK 8 有什么不同? JDK 7 用 Segment 分段锁(继承 ReentrantLock,默认 16 段,并发度被段数锁死);JDK 8 废弃 Segment,改成 Node 数组 + 链表/红黑树,空桶 CAS 插入、非空桶 synchronized 锁桶头;计数改 baseCount + CounterCell[];扩容支持多线程 helpTransfer 协助迁移,用 ForwardingNode(hash=MOVED)标记已迁移桶。
8.3 ThreadLocal Q13:ThreadLocal 为什么会发生内存泄漏?如何避免? Entry 的 key 是弱引用,GC 后变 null,但 value 仍是强引用,被 Thread → threadLocals → Entry → value 拽着;线程池核心线程长期存活,value 就永不回收。避免方式:用完在 finally 中 remove()。注意 expungeStaleEntry 的启发式清理是被动的,不能替代 remove()。
Q14:为什么 key 用弱引用而 value 用强引用? key 的生命周期由外部持有者决定(通常 static final),弱引用可在持有者消失后把 Entry 标成脏条目;value 的唯一强引用就来自 Entry,若设为弱引用,GC 后可能变 null,导致业务读到 null 的诡异 bug,比泄漏更难排查。故 value 必须强引用,靠 remove() 主动断开。
Q15:线程池里用 ThreadLocal 有什么风险? 一是脏数据 :线程复用,上次没清理的值被下次请求读到,可能造成跨用户数据泄露;二是内存泄漏 :线程不死,value 不回收。解法是在 Filter/Interceptor/AOP 统一设置与清理,remove() 放 finally;异步子线程用 TransmittableThreadLocal 透传上下文。
8.4 编码题 Q16:三个线程按顺序打印 ABC 各 10 次(用 ReentrantLock + Condition 实现)
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 import java.util.concurrent.locks.Condition;import java.util.concurrent.locks.ReentrantLock;public class PrintABC { private static final ReentrantLock lock = new ReentrantLock (); private static final Condition cA = lock.newCondition(); private static final Condition cB = lock.newCondition(); private static final Condition cC = lock.newCondition(); private static int state = 0 ; private static final int ROUND = 10 ; static class Printer implements Runnable { private final String name; private final int target; private final Condition self; private final Condition next; private final int nextState; Printer(String name, int target, Condition self, Condition next, int nextState) { this .name = name; this .target = target; this .self = self; this .next = next; this .nextState = nextState; } @Override public void run () { for (int i = 0 ; i < ROUND; i++) { lock.lock(); try { while (state != target) self.await(); System.out.print(name); state = nextState; next.signal(); } catch (InterruptedException e) { Thread.currentThread().interrupt(); return ; } finally { lock.unlock(); } } } } public static void main (String[] args) throws InterruptedException { Thread a = new Thread (new Printer ("A" , 0 , cA, cB, 1 ), "T-A" ); Thread b = new Thread (new Printer ("B" , 1 , cB, cC, 2 ), "T-B" ); Thread c = new Thread (new Printer ("C" , 2 , cC, cA, 0 ), "T-C" ); a.start(); b.start(); c.start(); a.join(); b.join(); c.join(); System.out.println("\n完成" ); } }
Q17:实现一个支持限流与超时回退的简易资源池
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60 61 62 63 64 65 66 67 68 import java.util.concurrent.*;import java.util.concurrent.atomic.AtomicInteger;public class SimpleConnectionPool { private final BlockingQueue<Connection> pool; private final Semaphore permits ; private final AtomicInteger created = new AtomicInteger (); static class Connection { final int id; Connection(int id) { this .id = id; } @Override public String toString () { return "Conn#" + id; } } public SimpleConnectionPool (int size) { this .pool = new ArrayBlockingQueue <>(size); this .permits = new Semaphore (size, true ); for (int i = 0 ; i < size; i++) pool.add(new Connection (i)); } public Connection borrow (long timeout, TimeUnit unit) throws InterruptedException { if (!permits .tryAcquire(timeout, unit)) throw new IllegalStateException ("连接池耗尽,等待超时" ); try { return pool.take(); } catch (InterruptedException e) { permits .release(); throw e; } } public void release (Connection c) { if (c == null ) return ; pool.offer(c); permits .release(); } public int available () { return permits .availablePermits(); } public static void main (String[] args) throws InterruptedException { SimpleConnectionPool p = new SimpleConnectionPool (3 ); ExecutorService es = Executors.newFixedThreadPool(10 ); CountDownLatch latch = new CountDownLatch (10 ); for (int i = 0 ; i < 10 ; i++) { es.submit(() -> { Connection c = null ; try { c = p.borrow(1 , TimeUnit.SECONDS); System.out.println(Thread.currentThread().getName() + " 借到 " + c); TimeUnit.MILLISECONDS.sleep(300 ); } catch (InterruptedException e) { Thread.currentThread().interrupt(); } catch (IllegalStateException e) { System.out.println(Thread.currentThread().getName() + " -> " + e.getMessage()); } finally { if (c != null ) p.release(c); latch.countDown(); } }); } latch.await(); System.out.println("归还后可用连接 = " + p.available()); es.shutdown(); } }
8.5 本篇小结 JUC 的心法浓缩成三句话:能不加锁就不加锁 (无锁 CAS、线程封闭、不可变对象);必须加锁就把粒度降到最小 (ConcurrentHashMap 的单个桶、LongAdder 的单个 Cell,本质都是”拆热点”);锁之外的协作交给现成工具 (Latch/Barrier/Semaphore/BlockingQueue)。下一篇进入 Executor 体系,看线程池如何把这些原语组装成工业级任务调度框架。