简介
在多线程程序设计中,有三个重要的同步工具需要我们掌握:Semaphore(信号量)、CountDownLatch(倒计数门闸锁)和 CyclicBarrier(可重用栅栏)。本文将重点介绍信号量 Semaphore 的原理、源码分析以及使用示例。
1. 信号量 Semaphore 的介绍
我们以一个停车场运作为例来说明信号量的作用。假设停车场只有三个车位,一开始三个车位都是空的。这时如果同时来了三辆车,看门人允许它们进入,然后放下车拦。以后来的车必须在入口等待,直到停车场中有车辆离开。这时,如果有一辆车离开停车场,看门人得知后,打开车拦,放入一辆;如果又离开一辆,则又可以放入一辆,如此往复。
在这个停车场系统中,车位是公共资源,每辆车好比一个线程,看门人起的就是信号量的作用。信号量是一个非负整数,表示当前公共资源的可用数目(在上面的例子中可以用空闲的停车位类比信号量)。当一个线程要使用公共资源时(在上面的例子中可以用车辆类比线程),首先要查看信号量:
- 如果信号量的值大于 0,则将其减 1,然后去占有公共资源。
- 如果信号量的值为 0,则线程会将自己阻塞,直到有其它线程释放公共资源。
在信号量上我们定义两种操作:acquire(获取)和release(释放)。
- 当一个线程调用 acquire 操作时,它要么成功获取信号量(信号量减 1),要么一直等待,直到有线程释放信号量或超时。
- release 操作会将信号量的值加 1,然后唤醒等待的线程。
信号量主要用于两个目的:
- 用于多个共享资源的互斥使用。
- 用于并发线程数的控制。
2. 信号量 Semaphore 的源码分析
在 Java 的并发包中,Semaphore类表示信号量。Semaphore 内部主要通过 AQS(AbstractQueuedSynchronizer)实现线程的管理。
2.1 构造函数
Semaphore 有两个构造函数,参数permits表示许可数,它最后传递给了 AQS 的state值。
// 非公平的构造函数 public Semaphore(int permits) { sync = new NonfairSync(permits); } // 通过 fair 参数决定公平性 public Semaphore(int permits, boolean fair) { sync = fair ? new FairSync(permits) : new NonfairSync(permits); }线程在运行时首先获取许可,如果成功,许可数就减 1,线程运行;当线程运行结束就释放许可,许可数就加 1。如果许可数为 0,则获取失败,线程位于 AQS 的等待队列中,它会被其它释放许可的线程唤醒。
在创建 Semaphore 对象的时候还可以指定它的公平性:
- 非公平信号量:在获取许可时先尝试获取许可,而不必关心是否已有需要获取许可的线程位于等待队列中,如果获取失败,才会入列。
- 公平信号量:在获取许可时首先要查看等待队列中是否已有线程,如果有则入列。
2.2 acquire 方法
public void acquire() throws InterruptedException { sync.acquireSharedInterruptibly(1); } public final void acquireSharedInterruptibly(int arg) throws InterruptedException { if (Thread.interrupted()) throw new InterruptedException(); if (tryAcquireShared(arg) < 0) doAcquireSharedInterruptibly(arg); } final int nonfairTryAcquireShared(int acquires) { for (;;) { int available = getState(); int remaining = available - acquires; if (remaining < 0 || compareAndSetState(available, remaining)) return remaining; } }可以看出,如果remaining < 0(即获取许可后,许可数小于 0),则获取失败,在doAcquireSharedInterruptibly方法中线程会将自身阻塞,然后入列。
2.3 release 方法
public void release() { sync.releaseShared(1); } public final boolean releaseShared(int arg) { if (tryReleaseShared(arg)) { doReleaseShared(); return true; } return false; } protected final boolean tryReleaseShared(int releases) { for (;;) { int current = getState(); int next = current + releases; if (next < current) // overflow throw new Error("Maximum permit count exceeded"); if (compareAndSetState(current, next)) return true; } }可以看出释放许可就是将 AQS 中state的值加 1,然后通过doReleaseShared唤醒等待队列的第一个节点。Semaphore 使用的是 AQS 的共享模式,等待队列中的第一个节点如果成功获取许可,又会唤醒下一个节点,以此类推。
3. 使用示例
下面是一个使用 Semaphore 控制并发线程数的示例:
package javalearning; import java.util.Random; import java.util.concurrent.ExecutorService; import java.util.concurrent.Executors; import java.util.concurrent.Semaphore; public class SemaphoreDemo { // 创建一个信号量,初始许可数为3,表示最多允许3个线程同时访问共享资源 private Semaphore smp = new Semaphore(3); private Random rnd = new Random(); class TaskDemo implements Runnable { private String id; TaskDemo(String id) { this.id = id; } @Override public void run() { try { // acquire() 方法:尝试获取一个许可 // 1. 如果当前有可用许可(state > 0),则获取成功,state减1,线程继续执行 // 2. 如果当前没有可用许可(state = 0),则线程会被阻塞,进入AQS等待队列 // 3. 该方法会响应中断,如果线程在等待时被中断,会抛出InterruptedException // 4. 注意:acquire() 是阻塞方法,会一直等待直到获取到许可或被中断 smp.acquire(); System.out.println("Thread " + id + " is working"); // 模拟线程执行任务,随机休眠0-999毫秒 Thread.sleep(rnd.nextInt(1000)); // release() 方法:释放一个许可 // 1. 将信号量的许可数加1(state加1) // 2. 如果有线程在等待队列中,会唤醒其中一个线程 // 3. 注意:release() 方法不会阻塞,总是立即返回 // 4. 重要:每个acquire()调用都应该有对应的release()调用,否则会导致许可泄漏 smp.release(); System.out.println("Thread " + id + " is over"); } catch (InterruptedException e) { // 异常处理块:处理线程中断异常 // 1. 当线程在acquire()等待时被中断,会进入这个catch块 // 2. 这里应该进行适当的清理工作,比如释放已获取的资源 // 3. 注意:在捕获InterruptedException后,通常需要恢复中断状态 // Thread.currentThread().interrupt(); // 4. 当前代码只是简单忽略中断,实际项目中应根据业务需求处理 // 5. 如果线程在sleep()时被中断,也会进入这个catch块 } } } public static void main(String[] args) { SemaphoreDemo semaphoreDemo = new SemaphoreDemo(); // 创建缓存线程池,会自动管理线程的创建和回收 ExecutorService se = Executors.newCachedThreadPool(); // 提交6个任务到线程池,但信号量只允许3个线程同时执行 se.submit(semaphoreDemo.new TaskDemo("a")); se.submit(semaphoreDemo.new TaskDemo("b")); se.submit(semaphoreDemo.new TaskDemo("c")); se.submit(semaphoreDemo.new TaskDemo("d")); se.submit(semaphoreDemo.new TaskDemo("e")); se.submit(semaphoreDemo.new TaskDemo("f")); // 关闭线程池,不再接受新任务,但会执行已提交的任务 se.shutdown(); // 注意:这里没有调用se.awaitTermination()等待所有任务完成 // 在实际应用中,可能需要等待所有任务完成后再结束程序 } }3.1 运行结果
Thread c is working Thread b is working Thread a is working Thread c is over Thread d is working Thread b is over Thread e is working Thread a is over Thread f is working Thread d is over Thread e is over Thread f is over可以看出,最多同时有三个线程并发执行,也可以认为有三个公共资源(比如计算机的三个串口)。
4. Semaphore 与 CountDownLatch、CyclicBarrier 的对比
在 Java 并发编程中,Semaphore、CountDownLatch 和 CyclicBarrier 都是重要的同步工具,它们各有不同的设计目的和使用场景。下面通过表格从多个维度进行对比:
| 对比维度 | Semaphore(信号量) | CountDownLatch(倒计数门闸锁) | CyclicBarrier(可重用栅栏) |
|---|---|---|---|
| 设计目的 | 控制同时访问特定资源的线程数量,管理有限资源的并发访问。 | 让一个或多个线程等待其他线程完成操作,实现线程间的协调。 | 让一组线程相互等待,直到所有线程都到达某个屏障点,然后同时继续执行。 |
| 核心机制 | 基于许可(permits)的计数器,acquire() 获取许可(计数器减1),release() 释放许可(计数器加1)。 | 基于倒计数的计数器,countDown() 减少计数,await() 等待计数归零。 | 基于屏障(barrier)的等待机制,await() 使线程等待,直到所有线程都调用了 await()。 |
| 可重用性 | 可重用,许可被释放后可被其他线程获取。 | 不可重用,计数归零后门闸打开,无法重置(除非新建实例)。 | 可重用,所有线程到达屏障后自动重置,可再次使用。 |
| 主要方法 | acquire()、release()、tryAcquire()、availablePermits() | await()、countDown()、getCount() | await()、reset()、getNumberWaiting()、getParties() |
| 典型场景 |
|
|
|
| 线程关系 | 通常用于限制资源访问的线程数量,线程之间是竞争关系。 | 一个或多个等待线程与一组工作线程之间的协调关系。 | 一组对等线程之间的相互等待关系。 |
| 计数器方向 | 许可数可增可减,acquire() 减少,release() 增加。 | 计数器只减不增,从初始值递减到0。 | 无计数器概念,基于到达屏障的线程数。 |
| 异常处理 | acquire() 可响应中断,release() 不会阻塞。 | await() 可响应中断和超时。 | await() 可响应中断、超时和屏障破坏异常。 |
总结与适用场景:
1.Semaphore最适合需要控制并发访问数量的场景。当你有有限数量的资源(如数据库连接、文件句柄、API调用配额)需要被多个线程共享时,Semaphore 可以确保同时访问资源的线程数不超过预设限制。它的核心思想是"资源配额管理"。
2.CountDownLatch最适合一次性协调场景。当你需要让一个或多个线程等待其他一组线程完成特定操作后才能继续执行时,CountDownLatch 是最佳选择。例如,主线程等待所有服务初始化完成,或者测试框架等待所有测试用例执行完毕。它的核心思想是"等待完成"。
3.CyclicBarrier最适合多阶段同步场景。当一组线程需要相互等待,在所有线程都到达某个点后才能继续执行下一阶段时,CyclicBarrier 非常有用。例如,并行计算中每轮迭代需要所有线程同步,或者多玩家游戏中每回合开始前等待所有玩家准备就绪。它的核心思想是"集体同步"。
选择建议:
- 需要限制资源访问数量 → 选择 Semaphore
- 需要等待其他线程完成 → 选择 CountDownLatch
- 需要线程组相互等待同步 → 选择 CyclicBarrier
在实际项目中,这三种工具可以结合使用,解决复杂的并发同步问题。理解它们的设计哲学和适用场景,有助于编写更高效、更安全的并发程序。
5. 信号量的典型应用场景
信号量在实际开发中有多种应用场景,以下是几个典型的例子:
5.1 数据库连接池
数据库连接是有限的资源,创建和维护连接需要消耗系统资源。使用信号量可以有效地管理数据库连接池:
- 初始化时创建固定数量的连接,并将信号量的许可数设置为连接数。
- 当线程需要获取数据库连接时,调用
acquire()方法获取许可。 - 如果连接池中有可用连接,线程立即获得连接并开始操作。
- 如果所有连接都被占用,线程会阻塞等待,直到有线程释放连接(调用
release())。 - 线程使用完连接后,必须调用
release()方法归还许可,以便其他线程可以使用。
这种方式可以防止过多的线程同时访问数据库,避免数据库过载,同时确保连接资源被高效复用。
5.2 限流器(Rate Limiter)
在高并发系统中,为了防止系统被突发流量冲垮,需要对请求进行限流:
- 设置信号量的许可数为系统能够承受的最大并发请求数。
- 每个请求到达时,首先尝试获取许可(
tryAcquire()或带超时的acquire())。 - 如果获取成功,请求被处理;如果获取失败(许可数为0),请求被拒绝或进入等待队列。
- 请求处理完成后释放许可,允许新的请求进入。
通过调整信号量的许可数,可以灵活控制系统的并发处理能力,保护后端服务不被过载。
5.3 生产者-消费者模型
在生产者-消费者模式中,信号量可以用于控制缓冲区的访问:
- 使用两个信号量:
emptySlots(空槽位信号量)和fullSlots(满槽位信号量)。 emptySlots初始值为缓冲区大小,表示可用空槽位数量。fullSlots初始值为0,表示已填充的槽位数量。- 生产者线程:先获取空槽位许可(
emptySlots.acquire()),生产数据放入缓冲区,然后释放满槽位许可(fullSlots.release())。 - 消费者线程:先获取满槽位许可(
fullSlots.acquire()),从缓冲区取出数据消费,然后释放空槽位许可(emptySlots.release())。
这种实现方式确保了生产者和消费者之间的同步,避免了缓冲区溢出或下溢的问题。
5.4 资源池管理
除了数据库连接池,信号量还可以用于管理其他类型的资源池:
- 线程池任务队列控制:限制同时等待执行的任务数量。
- 文件句柄管理:限制同时打开的文件数量,防止系统文件描述符耗尽。
- 网络连接限制:控制同时建立的网络连接数。
- 硬件设备访问:如打印机、扫描仪等共享设备的访问控制。
5.5 并发任务控制
在某些场景下,需要限制同时执行的特定类型任务数量:
- 批量数据处理:控制同时处理的数据分片数量,避免内存溢出。
- API调用限制:遵守第三方API的调用频率限制。
- 下载任务管理:限制同时进行的下载任务数量,避免网络拥堵。
信号量的灵活性和简单性使其成为并发编程中不可或缺的工具。通过合理设置许可数量和选择合适的获取/释放策略,可以解决多种并发控制问题。
6. 使用 Semaphore 的注意事项与最佳实践
虽然 Semaphore 是一个强大的并发控制工具,但在实际使用中需要注意一些常见的陷阱和最佳实践,以确保程序的正确性和性能。
6.1 许可泄漏(acquire 后未 release)
问题描述:许可泄漏是最常见的问题之一。当线程调用acquire()获取许可后,如果因为异常、逻辑错误或忘记调用release(),导致许可没有被释放,那么可用的许可数会逐渐减少,最终可能导致所有线程都无法获取许可而永久阻塞。
最佳实践:使用 try-finally 块确保release()一定会被调用。
public class SemaphoreSafeDemo { private Semaphore semaphore = new Semaphore(3); public void doWork() { try { semaphore.acquire(); // 执行业务逻辑 performTask(); } catch (InterruptedException e) { Thread.currentThread().interrupt(); // 恢复中断状态 // 处理中断逻辑 } finally { semaphore.release(); // 确保释放许可 } } private void performTask() { // 模拟业务逻辑 } }6.2 异常处理不当导致许可未释放
问题描述:当业务逻辑抛出未捕获的异常时,如果release()调用在异常之后,许可可能无法被释放。
最佳实践:将业务逻辑放在 try 块中,在 finally 块中释放许可。对于可中断的方法,要正确处理 InterruptedException。
public class ExceptionSafeDemo { private Semaphore semaphore = new Semaphore(2); public void processResource() { semaphore.acquireUninterruptibly(); // 使用不可中断的获取方式 try { // 可能抛出异常的业务逻辑 riskyOperation(); } finally { semaphore.release(); } } public void processWithTimeout() throws InterruptedException { if (semaphore.tryAcquire(1, TimeUnit.SECONDS)) { // 带超时的尝试获取 try { // 业务逻辑 doWork(); } finally { semaphore.release(); } } else { // 超时处理逻辑 handleTimeout(); } } private void riskyOperation() { // 可能抛出 RuntimeException 的操作 } private void doWork() { // 正常业务逻辑 } private void handleTimeout() { // 超时处理 } }6.3 tryAcquire 与 acquire 的选择策略
选择建议:
- 使用
acquire()的场景:当线程必须获取到许可才能继续执行时,使用阻塞式的acquire()。例如数据库连接池中,线程必须等待直到有可用连接。 - 使用
tryAcquire()的场景:当获取许可是可选的,或者需要快速失败时。例如限流器中,当系统过载时直接拒绝请求而不是让请求等待。 - 使用带超时的
tryAcquire(long timeout, TimeUnit unit)的场景:当需要限制等待时间,避免线程无限期阻塞时。
public class AcquisitionStrategyDemo { private Semaphore semaphore = new Semaphore(5); // 场景1:必须获取许可 - 使用 acquire() public void mustAcquire() throws InterruptedException { semaphore.acquire(); try { criticalOperation(); } finally { semaphore.release(); } } // 场景2:快速失败 - 使用 tryAcquire() public boolean tryAcquireFast() { if (semaphore.tryAcquire()) { try { optionalOperation(); return true; } finally { semaphore.release(); } } return false; // 立即返回,不阻塞 } // 场景3:有限等待 - 使用带超时的 tryAcquire() public boolean tryAcquireWithTimeout() throws InterruptedException { if (semaphore.tryAcquire(500, TimeUnit.MILLISECONDS)) { try { timeSensitiveOperation(); return true; } finally { semaphore.release(); } } return false; // 超时后返回 } private void criticalOperation() { // 必须执行的操作 } private void optionalOperation() { // 可选的操作 } private void timeSensitiveOperation() { // 对时间敏感的操作 } }6.4 公平与非公平模式对性能的影响
公平模式(FairSync):
- 特点:严格按照线程等待的先后顺序分配许可,先到先得。
- 优点:避免线程饥饿,保证公平性。
- 缺点:性能较低,因为需要维护等待队列,上下文切换开销大。
- 适用场景:当避免线程饥饿比性能更重要时,或者当等待时间可能很长时。
非公平模式(NonfairSync,默认):
- 特点:允许新请求的线程"插队",可能比等待队列中的线程先获取许可。
- 优点:性能更高,减少了线程切换的开销。
- 缺点:可能导致线程饥饿,某些线程可能长时间无法获取许可。
- 适用场景:大多数情况下的默认选择,特别是当许可持有时间很短或线程竞争不激烈时。
public class FairnessDemo { // 公平信号量 - 性能较低但保证公平 private Semaphore fairSemaphore = new Semaphore(3, true); // 非公平信号量(默认)- 性能较高但可能不公平 private Semaphore nonFairSemaphore = new Semaphore(3, false); // 或简写为:private Semaphore nonFairSemaphore = new Semaphore(3); public void testFairSemaphore() throws InterruptedException { fairSemaphore.acquire(); try { // 公平模式下,等待时间最长的线程优先获取许可 System.out.println("Fair semaphore acquired by: " + Thread.currentThread().getName()); Thread.sleep(100); } finally { fairSemaphore.release(); } } public void testNonFairSemaphore() throws InterruptedException { nonFairSemaphore.acquire(); try { // 非公平模式下,新请求的线程可能"插队" System.out.println("Non-fair semaphore acquired by: " + Thread.currentThread().getName()); Thread.sleep(100); } finally { nonFairSemaphore.release(); } } }6.5 其他注意事项
1. 避免在持有许可时执行耗时操作:
// 不推荐:在持有许可时执行耗时IO操作 semaphore.acquire(); try { // 耗时操作 - 这会长时间占用许可 processLargeFile(); // 可能执行几分钟 sendNetworkRequest(); // 网络延迟不可控 } finally { semaphore.release(); } // 推荐:尽快释放许可 List<String> data = prepareData(); // 准备数据(不持有许可) semaphore.acquire(); try { // 只执行必须同步的核心操作 updateSharedResource(data); } finally { semaphore.release(); // 尽快释放 } // 后续的非同步操作在释放许可后执行 logOperationResult();2. 合理设置许可数量:
- 设置过少:可能导致性能瓶颈,线程频繁等待。
- 设置过多:失去并发控制的意义,可能耗尽系统资源。
- 建议:根据系统资源(CPU核心数、内存、IO能力)和业务需求动态调整。
3. 监控信号量状态:
public class SemaphoreMonitor { private Semaphore semaphore = new Semaphore(10); public void printStatus() { System.out.println("可用许可数: " + semaphore.availablePermits()); System.out.println("等待队列长度(估算): " + semaphore.getQueueLength()); System.out.println("是否有线程正在等待: " + semaphore.hasQueuedThreads()); } public void acquireWithLogging() throws InterruptedException { System.out.println("尝试获取许可,当前可用: " + semaphore.availablePermits()); semaphore.acquire(); System.out.println("成功获取许可,剩余可用: " + semaphore.availablePermits()); } }4. 避免嵌套获取:
// 危险:可能导致死锁 semaphore.acquire(); try { // 某些条件下再次获取同一个信号量 if (needMoreResource()) { semaphore.acquire(); // 如果许可数不足,这里会永久阻塞! try { // ... } finally { semaphore.release(); } } } finally { semaphore.release(); } // 解决方案:使用可重入锁或其他同步机制 private ReentrantLock lock = new ReentrantLock(); public void safeNestedAccess() { lock.lock(); try { if (needMoreResource()) { // 使用其他同步机制而不是嵌套获取同一个信号量 handleAdditionalResource(); } } finally { lock.unlock(); } }通过遵循这些最佳实践,可以避免 Semaphore 使用中的常见陷阱,编写出更健壮、高效的并发程序。
7. 总结
Semaphore 是 Java 并发编程中重要的同步工具,通过控制许可数量来管理对共享资源的访问。本文通过停车场示例介绍了信号量的基本概念,分析了 Semaphore 的源码实现,并提供了完整的使用示例。掌握 Semaphore 的使用可以帮助我们更好地设计并发程序。