Semaphore控制并发访问数量

wen java案例 1

深入解析 Semaphore:高效控制并发访问数量的终极指南

📖 目录导读

  1. 什么是 Semaphore?—— 并发控制的核心概念
  2. Semaphore 的工作原理与信号量机制详解
  3. 实际应用场景:从数据库连接到限流器
  4. 代码实战:Java/Python/Go 中的 Semaphore 实现
  5. 常见错误与性能陷阱:如何避免死锁与饥饿
  6. 与其它并发工具的对比:Semaphore vs 互斥锁 vs 线程池
  7. 问答环节:解决你最关心的 5 个问题
  8. 掌握 Semaphore 的黄金法则

什么是 Semaphore?—— 并发控制的核心概念

在现代软件开发中,并发访问控制 是保证系统稳定性的关键,想象一下:你有一个只能同时容纳 10 人的电梯,100 人同时涌进去,后果不堪设想。Semaphore(信号量) 就是这样一个“数字门禁”——它通过一个计数器来限制同时访问特定资源的线程/进程数量。

Semaphore控制并发访问数量

定义:Semaphore 是一个非负整数计数器,用于控制对共享资源的访问,当计数器大于 0 时,线程可以获取许可(acquire)并继续执行;当计数器为 0 时,线程必须等待,直到其他线程释放(release)许可。

为什么需要 Semaphore?

  • 限制并发请求:防止数据库连接池被耗尽
  • 流量整形:平滑处理突发流量,如秒杀系统
  • 资源保护:确保有限资源(如打印机、GPU 实例)不被过度使用

Semaphore 的工作原理与信号量机制详解

核心操作

操作 说明 伪代码
acquire() 请求一个许可;若计数器 >0,则减 1 并继续;否则阻塞等待 while (count == 0) { wait; } count--;
release() 释放一个许可,计数器加 1,并唤醒等待线程 count++; notify();
tryAcquire() 非阻塞尝试获取许可,立即返回 true/false if (count > 0) { count--; return true; } else return false;

两种主要类型

  • 计数信号量(Counting Semaphore):许可数 >1,用于控制资源池(如连接池)
  • 二元信号量(Binary Semaphore):许可数 =1,等价于互斥锁(Mutex)

底层实现机制

操作系统通过 PV 操作(P 表示 acquire,V 表示 release)实现原子性增减,在 Java 中,Semaphore 类基于 AQS(AbstractQueuedSynchronizer)框架,能够公平模式(FIFO 队列)或非公平模式(抢占式)工作。

关键特性:Semaphore 不会绑定到特定线程——任何线程都可以释放许可,这比互斥锁更灵活。


实际应用场景:从数据库连接到限流器

数据库连接池限流

Semaphore pool = new Semaphore(20); // 最多20个连接
// 请求连接
pool.acquire();
Connection conn = getConnection();
try {
    // 执行查询
} finally {
    releaseConnection(conn);
    pool.release(); // 释放许可
}

API 流量控制(Guava RateLimiter 的替代方案)

import threading
semaphore = threading.Semaphore(5)  # 每秒最多5个请求
def handle_request():
    with semaphore:
        # 处理请求
        pass

多线程任务分发

var sem = make(chan struct{}, 3) // 最多3个并发goroutine
for _, task := range tasks {
    sem <- struct{}{} // acquire
    go func(t Task) {
        defer func() { <-sem }() // release
        t.execute()
    }(task)
}

更复杂的场景:文件下载限速

结合计数器实现动态限流,如限制各 IP 的并发下载数。


代码实战:Java/Python/Go 中的 Semaphore 实现

Java 示例(公平模式)

import java.util.concurrent.Semaphore;
public class PrintQueue {
    private final Semaphore semaphore = new Semaphore(3, true); // 公平模式
    public void printJob(Object document) {
        try {
            semaphore.acquire();
            System.out.println(Thread.currentThread().getName() + " 开始打印");
            Thread.sleep(2000);
            System.out.println(Thread.currentThread().getName() + " 打印完毕");
        } catch (InterruptedException e) {
            Thread.currentThread().interrupt();
        } finally {
            semaphore.release();
        }
    }
}

Python 示例(async+await 异步版本)

import asyncio
async def worker(sem, name):
    async with sem:
        print(f'{name} 获得许可')
        await asyncio.sleep(1)
    print(f'{name} 释放许可')
async def main():
    sem = asyncio.Semaphore(2)
    tasks = [worker(sem, f'Worker-{i}') for i in range(5)]
    await asyncio.gather(*tasks)
asyncio.run(main())

Go 示例(channel 实现信号量)

package main
import (
    "fmt"
    "sync"
    "time"
)
type Semaphore struct {
    ch chan struct{}
}
func NewSemaphore(max int) *Semaphore {
    return &Semaphore{ch: make(chan struct{}, max)}
}
func (s *Semaphore) Acquire() {
    s.ch <- struct{}{}
}
func (s *Semaphore) Release() {
    <-s.ch
}
func main() {
    sem := NewSemaphore(2)
    var wg sync.WaitGroup
    for i := 0; i < 5; i++ {
        wg.Add(1)
        go func(id int) {
            defer wg.Done()
            sem.Acquire()
            defer sem.Release()
            fmt.Printf("Worker %d 开始\n", id)
            time.Sleep(1 * time.Second)
            fmt.Printf("Worker %d 结束\n", id)
        }(i)
    }
    wg.Wait()
}

常见错误与性能陷阱:如何避免死锁与饥饿

错误 1:忘记释放许可(Leak)

// 错误示例:如果异常发生,release 不会执行
semaphore.acquire();
doSth(); // 可能抛出异常
semaphore.release(); // 不执行!

解决方案:始终在 finally 块中释放,或使用 try-with-resources(Java 7+)

错误 2:许可数设置不合理

  • 太小:导致大量线程阻塞,CPU 浪费在上下文切换
  • 太大:等于没限制,资源依然耗竭
  • 黄金法则许可数 = 资源物理限制 × 安全系数(0.7-0.8)

错误 3:死锁与线程阻塞

当多个 Semaphore 嵌套使用时,可能出现:

  • 线程 A 持有信号量 1,等待信号量 2
  • 线程 B 持有信号量 2,等待信号量 1

解决:统一获取顺序,或使用 tryAcquire 带超时

性能陷阱:公平模式 vs 非公平模式

  • 非公平模式(默认):吞吐量高,但可能线程饥饿
  • 公平模式:避免饥饿,但性能稍低(需要维护 FIFO 队列)
  • 选择建议:长任务用公平模式,短任务用非公平模式

与其它并发工具的对比

特性 Semaphore Mutex(互斥锁) 线程池 CountDownLatch
目标 控制并发数 互斥访问 管理线程生命周期 等待多个线程完成
许可数 可多个(N) 仅1个 取决于池大小 不可复用
释放者 任意线程 只能持有锁的线程 N/A 所有线程
适用场景 限流、资源池 临界区保护 任务队列执行 初始化等待

何时选择 Semaphore 而非线程池?

  • 线程池控制的是执行线程数,Semaphore 控制的是资源访问数
  • 数据库连接池(资源数固定)用 Semaphore;计算密集型任务(线程数影响 CPU)用线程池

问答环节:解决你最关心的 5 个问题

❓ Q1:Semaphore 可以动态调整许可数吗?

A:标准库不支持直接动态调整,但可以通过 drainPermits() 清空后再 release() 变相修改。

semaphore.drainPermits(); // 清空所有许可
semaphore.release(10); // 重新设置为10

❓ Q2:如何实现限流器(每秒 X 个请求)?

A:结合时间窗口和 Semaphore:

ScheduledExecutorService scheduler = Executors.newScheduledThreadPool(1);
final Semaphore limiter = new Semaphore(10); // 每秒10个
scheduler.scheduleAtFixedRate(() -> limiter.release(10 - limiter.availablePermits()), 1, 1, TimeUnit.SECONDS);

(注意:更推荐使用 Guava RateLimiter)

❓ Q3:Semaphore 与 Redis 限流的区别?

A:Semaphore 是应用层本地限流,适合单机;Redis 限流(如令牌桶)适合分布式场景,但网络开销更大。

❓ Q4:acquire 时被中断怎么办?

A:正确做法是 catch (InterruptedException) 并调用 Thread.currentThread().interrupt() 恢复中断状态,然后释放已获取的资源。

❓ Q5:为什么推荐 tryAcquire(timeout) 而不是 acquire

Aacquire() 会永久阻塞,可能导致系统挂起;tryAcquire(5, TimeUnit.SECONDS) 给系统留有余地,避免死锁累积。


掌握 Semaphore 的黄金法则

三要三不要

✅ 应该做 ❌ 不要做
在 finally 块中释放许可 依赖 GC 自动回收
合理评估许可数(压力测试) 随意设置过大或过小
结合超时+重试机制 无限制阻塞等待

核心记忆口诀

获取务必释放,许可数量科学;公平非公权衡,超时避免死锁。

请记住:Semaphore 是工具,不是银弹,当你的系统需要精细控制并发访问数量时,它是性价比最高的选择;但如果你的目标是保护临界区(如写入共享变量),请使用 ReentrantLocksynchronized

你已经掌握了信号量的精髓——去优化你的并发系统吧!

抱歉,评论功能暂时关闭!