返回文章列表
JUC并发编程
JUC读写锁SemaphoreCountDownLatchCyclicBarrier

09JUC锁与同步工具

JUC(java.util.concurrent)中除了线程池,还提供了大量锁与同步工具。本篇先对 AQS、ReentrantLock 做概念性概述,再重点学习读写锁 ReentrantReadWriteLock、乐观读 StampedLock,以及三个常用同步工具:Semaphore、CountDownLatch、CyclicBarrier。

AQS 原理概述

AQS 全称 AbstractQueuedSynchronizer(抽象队列同步器),是 JUC 中几乎所有锁和同步工具的共同基础:

  • ReentrantLock
  • ReentrantReadWriteLock
  • Semaphore
  • CountDownLatch
  • CyclicBarrier(基于 ReentrantLock + Condition 间接使用 AQS)
  • 线程池中的 Worker 等

AQS 的核心思想可以概括为三点:

  1. state 状态变量:用一个 volatile int state 表示同步状态。对 ReentrantLock 而言 state 是加锁次数(0 表示未加锁);对 Semaphore 而言是剩余许可数;对 CountDownLatch 而言是倒计时计数
  2. CAS 修改状态:多线程通过 CAS 操作竞争修改 state,保证修改的原子性
  3. 阻塞队列:竞争失败的线程被封装成节点,加入一个 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();
        }
    }
}

注意两个细节:

  1. 读锁发现缓存失效后,必须先释放读锁再获取写锁(读写锁不支持升级)
  2. 获取到写锁后要再次检查 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 一组线程互相等待,计数满足后一起执行 是,可循环使用