09JUC锁与同步工具
JUC(java.util.concurrent)中除了线程池,还提供了大量锁与同步工具。本篇先对 AQS、ReentrantLock 做概念性概述,再重点学习读写锁 ReentrantReadWriteLock、乐观读 StampedLock,以及三个常用同步工具:Semaphore、CountDownLatch、CyclicBarrier。
AQS 原理概述
AQS 全称 AbstractQueuedSynchronizer(抽象队列同步器),是 JUC 中几乎所有锁和同步工具的共同基础:
ReentrantLockReentrantReadWriteLockSemaphoreCountDownLatchCyclicBarrier(基于ReentrantLock+Condition间接使用 AQS)- 线程池中的 Worker 等
AQS 的核心思想可以概括为三点:
- state 状态变量:用一个
volatile int state表示同步状态。对 ReentrantLock 而言 state 是加锁次数(0 表示未加锁);对 Semaphore 而言是剩余许可数;对 CountDownLatch 而言是倒计时计数 - CAS 修改状态:多线程通过 CAS 操作竞争修改 state,保证修改的原子性
- 阻塞队列:竞争失败的线程被封装成节点,加入一个 FIFO 的双向等待队列(CLH 队列的变体)排队挂起,状态释放后再唤醒队列中的线程
此外 AQS 还分独占模式(同一时刻只有一个线程能获取,如 ReentrantLock)和共享模式(多个线程可以同时获取,如 Semaphore、CountDownLatch)。
AQS 与 ReentrantLock 的深入源码(加解锁流程、队列迁移、公平与非公平实现)在原理篇另讲,本篇只做概念性了解。
ReentrantLock 原理概述
ReentrantLock 是基于 AQS 独占模式实现的可重入锁,相比 synchronized 它提供了更灵活的能力:
- 可重入:一个线程获取锁之后,可以再次重复获取同一把锁,内部用 state 记录重入次数,释放时也要释放相同次数
- 可中断:
lockInterruptibly()支持在等待锁的过程中被中断 - 可超时:
tryLock(timeout, unit)等待指定时间后仍获取不到锁就放弃 - 支持公平锁与非公平锁:构造时传入
true为公平锁(按排队顺序获取锁),默认false为非公平锁(允许插队,吞吐量更高) - 支持多个条件变量:可以创建多个
Condition,实现比wait/notify更精细的等待/通知
使用范式:
ReentrantLock lock = new ReentrantLock();
lock.lock();
try {
// 临界区
} finally {
lock.unlock(); // 必须在 finally 中释放
}读写锁 ReentrantReadWriteLock
当读操作远远高于写操作时,使用读写锁让读-读可以并发,提高性能,类似于数据库中的 select ... from ... lock in share mode(共享锁)。互斥规则为:
- 读-读:不互斥,可以并发
- 读-写:互斥
- 写-写:互斥
基本使用
提供一个数据容器类,内部用读锁保护 read() 方法,写锁保护 write() 方法:
class DataContainer {
private Object data;
private ReentrantReadWriteLock rw = new ReentrantReadWriteLock();
private ReentrantReadWriteLock.ReadLock r = rw.readLock();
private ReentrantReadWriteLock.WriteLock w = rw.writeLock();
public Object read() {
log.debug("获取读锁...");
r.lock();
try {
log.debug("读取");
sleep(1);
return data;
} finally {
log.debug("释放读锁...");
r.unlock();
}
}
public void write() {
log.debug("获取写锁...");
w.lock();
try {
log.debug("写入");
sleep(1);
} finally {
log.debug("释放写锁...");
w.unlock();
}
}
}测试读锁-读锁可以并发:
DataContainer dataContainer = new DataContainer();
new Thread(() -> {
dataContainer.read();
}, "t1").start();
new Thread(() -> {
dataContainer.read();
}, "t2").start();输出结果可以看到,t1 持有读锁期间,t2 的读操作不受影响,两个线程同时进入读取:
14:05:14.341 c.DataContainer [t2] - 获取读锁...
14:05:14.341 c.DataContainer [t1] - 获取读锁...
14:05:14.345 c.DataContainer [t1] - 读取
14:05:14.345 c.DataContainer [t2] - 读取
14:05:15.365 c.DataContainer [t2] - 释放读锁...
14:05:15.386 c.DataContainer [t1] - 释放读锁...测试读锁-写锁相互阻塞:
DataContainer dataContainer = new DataContainer();
new Thread(() -> {
dataContainer.read();
}, "t1").start();
Thread.sleep(100);
new Thread(() -> {
dataContainer.write();
}, "t2").start();输出结果,写操作必须等读锁释放后才能执行:
14:04:21.838 c.DataContainer [t1] - 获取读锁...
14:04:21.838 c.DataContainer [t2] - 获取写锁...
14:04:21.841 c.DataContainer [t2] - 写入
14:04:22.843 c.DataContainer [t2] - 释放写锁...
14:04:22.843 c.DataContainer [t1] - 读取
14:04:23.843 c.DataContainer [t1] - 释放读锁...写锁-写锁也是相互阻塞的,此处不再测试。
注意事项
- 读锁不支持条件变量
- 重入时升级不支持:持有读锁的情况下去获取写锁,会导致永久等待
r.lock();
try {
// ...
w.lock(); // 持有读锁再获取写锁 → 永久等待,不要这样写
try {
// ...
} finally {
w.unlock();
}
} finally {
r.unlock();
}- 重入时降级支持:持有写锁的情况下可以再获取读锁
锁降级与缓存一致性应用
读写锁的一个典型应用场景是缓存。缓存有效时直接读(多个线程可并发读);缓存失效时,只有一个线程获取写锁去重新计算数据,更新完成后降级为读锁再释放写锁,保证当前线程看到的数据一致、同时让其它线程也能读取缓存:
class CachedData {
Object data;
// 是否有效,如果失效,需要重新计算 data
volatile boolean cacheValid;
final ReentrantReadWriteLock rwl = new ReentrantReadWriteLock();
void processCachedData() {
rwl.readLock().lock();
if (!cacheValid) {
// 获取写锁前必须释放读锁
rwl.readLock().unlock();
rwl.writeLock().lock();
try {
// 判断是否有其它线程已经获取了写锁、更新了缓存, 避免重复更新(双重检查)
if (!cacheValid) {
data = ...; // 重新计算
cacheValid = true;
}
// 降级为读锁, 释放写锁, 这样能够让其它线程读取缓存
rwl.readLock().lock();
} finally {
rwl.writeLock().unlock();
}
}
// 自己用完数据, 释放读锁
try {
use(data);
} finally {
rwl.readLock().unlock();
}
}
}注意两个细节:
- 读锁发现缓存失效后,必须先释放读锁再获取写锁(读写锁不支持升级)
- 获取到写锁后要再次检查
cacheValid(双重检查),避免多个排队线程重复更新缓存
StampedLock
StampedLock 自 JDK 8 加入,是为了进一步优化读性能。它的特点是在使用读锁、写锁时都必须配合【戳(stamp)】使用。
加解读锁:
long stamp = lock.readLock();
lock.unlockRead(stamp);加解写锁:
long stamp = lock.writeLock();
lock.unlockWrite(stamp);乐观读:StampedLock 支持 tryOptimisticRead() 方法。读取完毕后需要做一次戳校验,如果校验通过,表示读取期间确实没有写操作,数据可以安全使用;如果校验没通过,需要重新获取读锁,保证数据安全:
long stamp = lock.tryOptimisticRead();
// 验戳
if (!lock.validate(stamp)) {
// 校验失败,升级为读锁
}数据容器示例:
class DataContainerStamped {
private int data;
private final StampedLock lock = new StampedLock();
public DataContainerStamped(int data) {
this.data = data;
}
public int read(int readTime) {
long stamp = lock.tryOptimisticRead();
log.debug("optimistic read locking...{}", stamp);
sleep(readTime);
if (lock.validate(stamp)) {
log.debug("read finish...{}, data:{}", stamp, data);
return data;
}
// 锁升级 - 读锁
log.debug("updating to read lock... {}", stamp);
try {
stamp = lock.readLock();
log.debug("read lock {}", stamp);
sleep(readTime);
log.debug("read finish...{}, data:{}", stamp, data);
return data;
} finally {
log.debug("read unlock {}", stamp);
lock.unlockRead(stamp);
}
}
public void write(int newData) {
long stamp = lock.writeLock();
log.debug("write lock {}", stamp);
try {
sleep(2);
this.data = newData;
} finally {
log.debug("write unlock {}", stamp);
lock.unlockWrite(stamp);
}
}
}测试读-读并发,可以看到实际没有加读锁(两个线程都走乐观读):
public static void main(String[] args) {
DataContainerStamped dataContainer = new DataContainerStamped(1);
new Thread(() -> {
dataContainer.read(1);
}, "t1").start();
sleep(0.5);
new Thread(() -> {
dataContainer.read(0);
}, "t2").start();
}输出:
15:58:50.217 c.DataContainerStamped [t1] - optimistic read locking...256
15:58:50.717 c.DataContainerStamped [t2] - optimistic read locking...256
15:58:50.717 c.DataContainerStamped [t2] - read finish...256, data:1
15:58:51.220 c.DataContainerStamped [t1] - read finish...256, data:1测试读-写并发,乐观读校验失败后补加读锁,保证读到的是写入后的新值:
public static void main(String[] args) {
DataContainerStamped dataContainer = new DataContainerStamped(1);
new Thread(() -> {
dataContainer.read(1);
}, "t1").start();
sleep(0.5);
new Thread(() -> {
dataContainer.write(100);
}, "t2").start();
}输出:
15:57:00.219 c.DataContainerStamped [t1] - optimistic read locking...256
15:57:00.717 c.DataContainerStamped [t2] - write lock 384
15:57:01.225 c.DataContainerStamped [t1] - updating to read lock... 256
15:57:02.719 c.DataContainerStamped [t2] - write unlock 384
15:57:02.719 c.DataContainerStamped [t1] - read lock 513
15:57:03.719 c.DataContainerStamped [t1] - read finish...513, data:1000
15:57:03.719 c.DataContainerStamped [t1] - read unlock 513注意:
- StampedLock 不支持条件变量
- StampedLock 不支持可重入
- 乐观读期间数据可能被修改,一定不能省略
validate(stamp)校验
Semaphore 信号量
Semaphore([ˈsɛməˌfɔr])信号量,用来限制能同时访问共享资源的线程上限,典型用途是限流。
基本使用:
public static void main(String[] args) {
// 1. 创建 semaphore 对象,许可数量为 3
Semaphore semaphore = new Semaphore(3);
// 2. 10个线程同时运行
for (int i = 0; i < 10; i++) {
new Thread(() -> {
// 3. 获取许可
try {
semaphore.acquire();
} catch (InterruptedException e) {
e.printStackTrace();
}
try {
log.debug("running...");
sleep(1);
log.debug("end...");
} finally {
// 4. 释放许可
semaphore.release();
}
}).start();
}
}输出中可以清楚看到,10 个线程每批最多只有 3 个并发执行,有线程释放许可后,等待的线程才能获取许可继续:
07:35:15.485 c.TestSemaphore [Thread-2] - running...
07:35:15.485 c.TestSemaphore [Thread-1] - running...
07:35:15.485 c.TestSemaphore [Thread-0] - running...
07:35:16.490 c.TestSemaphore [Thread-2] - end...
07:35:16.490 c.TestSemaphore [Thread-0] - end...
07:35:16.490 c.TestSemaphore [Thread-1] - end...
07:35:16.490 c.TestSemaphore [Thread-3] - running...
07:35:16.490 c.TestSemaphore [Thread-5] - running...
07:35:16.490 c.TestSemaphore [Thread-4] - running...
...常用 API:
acquire()/acquire(n):获取 1 个 / n 个许可,没有许可时阻塞tryAcquire()/tryAcquire(timeout, unit):尝试获取许可,获取不到立即返回 false 或等待超时release()/release(n):释放 1 个 / n 个许可- 构造时传入
true可使用公平模式(按等待顺序发放许可)
可以把 Semaphore 理解为「停车场入口」:许可数就是车位数,车位满后新来的车(线程)必须等待,有车离开(release)后才能进入。
CountDownLatch
CountDownLatch 用来进行线程同步协作,等待所有线程完成倒计时。
- 构造参数用来初始化等待计数值
await()用来等待计数归零countDown()用来让计数减一
基本使用
public static void main(String[] args) throws InterruptedException {
CountDownLatch latch = new CountDownLatch(3);
new Thread(() -> {
log.debug("begin...");
sleep(1);
latch.countDown();
log.debug("end...{}", latch.getCount());
}).start();
new Thread(() -> {
log.debug("begin...");
sleep(2);
latch.countDown();
log.debug("end...{}", latch.getCount());
}).start();
new Thread(() -> {
log.debug("begin...");
sleep(1.5);
latch.countDown();
log.debug("end...{}", latch.getCount());
}).start();
log.debug("waiting...");
latch.await();
log.debug("wait end...");
}输出:
18:44:00.778 c.TestCountDownLatch [main] - waiting...
18:44:00.778 c.TestCountDownLatch [Thread-2] - begin...
18:44:00.778 c.TestCountDownLatch [Thread-0] - begin...
18:44:00.778 c.TestCountDownLatch [Thread-1] - begin...
18:44:01.782 c.TestCountDownLatch [Thread-0] - end...2
18:44:02.283 c.TestCountDownLatch [Thread-2] - end...1
18:44:02.782 c.TestCountDownLatch [Thread-1] - end...0
18:44:02.782 c.TestCountDownLatch [main] - wait end...配合线程池使用
public static void main(String[] args) throws InterruptedException {
CountDownLatch latch = new CountDownLatch(3);
ExecutorService service = Executors.newFixedThreadPool(4);
service.submit(() -> {
log.debug("begin...");
sleep(1);
latch.countDown();
log.debug("end...{}", latch.getCount());
});
service.submit(() -> {
log.debug("begin...");
sleep(1.5);
latch.countDown();
log.debug("end...{}", latch.getCount());
});
service.submit(() -> {
log.debug("begin...");
sleep(2);
latch.countDown();
log.debug("end...{}", latch.getCount());
});
service.submit(() -> {
try {
log.debug("waiting...");
latch.await();
log.debug("wait end...");
} catch (InterruptedException e) {
e.printStackTrace();
}
});
}典型应用
- 同步等待多线程准备完毕:例如 10 个玩家线程各自加载资源(进度 0%~100%),主线程用
latch.await()等待所有玩家加载完成后才打印「游戏开始」 - 同步等待多个远程调用结束:一个接口需要同时调用订单、商品、物流等多个远程服务,可以用线程池并发发起调用,每个调用结束后
countDown(),主线程等待所有调用结束后再汇总结果。相比串行调用,总耗时从各调用耗时之和缩短为最慢的那个调用的耗时,伪代码如下:
RestTemplate restTemplate = new RestTemplate();
log.debug("begin");
ExecutorService service = Executors.newCachedThreadPool();
CountDownLatch latch = new CountDownLatch(4);
Future<Map<String, Object>> f1 = service.submit(() ->
restTemplate.getForObject("http://localhost:8080/order/{1}", Map.class, 1));
Future<Map<String, Object>> f2 = service.submit(() ->
restTemplate.getForObject("http://localhost:8080/product/{1}", Map.class, 1));
Future<Map<String, Object>> f3 = service.submit(() ->
restTemplate.getForObject("http://localhost:8080/product/{1}", Map.class, 2));
Future<Map<String, Object>> f4 = service.submit(() ->
restTemplate.getForObject("http://localhost:8080/logistics/{1}", Map.class, 1));
System.out.println(f1.get());
System.out.println(f2.get());
System.out.println(f3.get());
System.out.println(f4.get());
log.debug("执行完毕");
service.shutdown();另外,await(timeout, unit) 还支持超时等待,避免个别线程故障导致主线程永久阻塞。
CyclicBarrier
CyclicBarrier([ˈsaɪklɪk ˈbæriɚ])循环栅栏,用来进行线程协作:等待线程满足某个计数。构造时设置「计数个数」,每个线程执行到需要同步的时刻调用 await() 等待,当等待的线程数满足计数个数时,所有线程再一起继续执行。
CyclicBarrier cb = new CyclicBarrier(2); // 个数为2时才会继续执行
new Thread(() -> {
System.out.println("线程1开始.." + new Date());
try {
cb.await(); // 当个数不足时,等待
} catch (InterruptedException | BrokenBarrierException e) {
e.printStackTrace();
}
System.out.println("线程1继续向下运行..." + new Date());
}).start();
new Thread(() -> {
System.out.println("线程2开始.." + new Date());
try { Thread.sleep(2000); } catch (InterruptedException e) { }
try {
cb.await(); // 2 秒后,线程个数够2,继续运行
} catch (InterruptedException | BrokenBarrierException e) {
e.printStackTrace();
}
System.out.println("线程2继续向下运行..." + new Date());
}).start();CyclicBarrier 还提供一个构造参数,可以指定计数满足后优先执行的回调任务:
CyclicBarrier cb = new CyclicBarrier(5, () -> {
System.out.println("人齐了,发车!");
});与 CountDownLatch 的区别
- CountDownLatch 的计数不能重置,是一次性的;CyclicBarrier 计数满足后可以重置并重复使用,所以叫「循环」栅栏
- CountDownLatch 是一个线程(或多个线程)等待另外若干个线程完成
countDown();CyclicBarrier 是一组线程互相等待,大家都到齐了再一起走 - CyclicBarrier 可以被比喻为「人满发车」:旅游大巴固定载客数,到了人数就发车,发完一班还可以继续等下一拨人
小结
| 工具 | 作用 | 是否可重用 |
|---|---|---|
| ReentrantReadWriteLock | 读读并发、读写/写写互斥,适合读多写少 | 是 |
| StampedLock | 在读写锁基础上提供乐观读,进一步提升读性能 | 是,但不可重入 |
| Semaphore | 限制同时访问资源的线程数(限流) | 许可可循环获取/释放 |
| CountDownLatch | 一个线程等待多个线程倒计时归零 | 否,一次性 |
| CyclicBarrier | 一组线程互相等待,计数满足后一起执行 | 是,可循环使用 |