一、信号量是什么
信号量是一个计数器,控制同时访问某个资源的线程/进程数量。
锁(lock):同一时刻只允许 1 个线程进入 → 互斥(二元信号量) 信号量(semaphore):同一时刻允许 n 个线程进入 → 限流/资源池
生活中的类比:
| 场景 | 信号量值 | 含义 |
|---|---|---|
| 停车场 | 100 | 最多 100 辆车同时停 |
| 餐厅座位 | 50 | 最多 50 人同时就餐 |
| 卫生间隔间 | 3 | 最多 3 人同时使用 |
| 电梯载重 | 10 | 最多 10 人同时乘坐 |
二、核心操作
信号量只有两个基本操作:
acquire() — 获取一个许可(计数器 -1) release() — 释放一个许可(计数器 +1) 初始许可数 = 3 线程a acquire → 剩余许可: 2 线程b acquire → 剩余许可: 1 线程c acquire → 剩余许可: 0 线程d acquire → 阻塞等待(许可为0,没有空位了) ... 线程a release → 剩余许可: 1 线程d 被唤醒 → 获取许可成功,剩余许可: 0
三、jdk 中的 semaphore
基本用法
import java.util.concurrent.semaphore;
// 创建信号量:最多允许 3 个线程同时执行
semaphore semaphore = new semaphore(3);
public void accessresource() {
try {
semaphore.acquire(); // 获取许可(阻塞等待)
// 临界区:最多 3 个线程同时在这里
dowork();
} catch (interruptedexception e) {
thread.currentthread().interrupt();
} finally {
semaphore.release(); // 释放许可
}
}
构造函数
// 非公平信号量(默认):不保证等待顺序 semaphore semaphore = new semaphore(3); // 公平信号量:按请求顺序获取许可(fifo) semaphore fairsemaphore = new semaphore(3, true);
常用方法
// 阻塞获取 1 个许可 semaphore.acquire(); // 阻塞获取多个许可 semaphore.acquire(2); // 一次获取 2 个 // 尝试获取,获取不到立即返回 false(不阻塞) boolean acquired = semaphore.tryacquire(); // 尝试获取,最多等待指定时间 boolean acquired = semaphore.tryacquire(5, timeunit.seconds); // 释放许可 semaphore.release(); // 查看当前可用许可数 int available = semaphore.availablepermits(); // 获取正在等待的线程数 int waiting = semaphore.getqueuelength();
示例:数据库连接池
public class simpleconnectionpool {
private final semaphore semaphore;
private final queue<connection> pool;
public simpleconnectionpool(int maxsize) {
this.semaphore = new semaphore(maxsize);
this.pool = new concurrentlinkedqueue<>();
// 预创建连接
for (int i = 0; i < maxsize; i++) {
pool.offer(createconnection());
}
}
public connection getconnection() throws interruptedexception {
semaphore.acquire(); // 获取许可(控制并发数)
return pool.poll(); // 取出连接
}
public void releaseconnection(connection conn) {
pool.offer(conn); // 归还连接
semaphore.release(); // 释放许可
}
}
示例:接口限流(本地)
@restcontroller
public class ordercontroller {
// 最多允许 10 个请求同时处理下单
private final semaphore ordersemaphore = new semaphore(10);
@postmapping("/order/create")
public result createorder(@requestbody orderdto dto) {
if (!ordersemaphore.tryacquire()) {
return result.fail("系统繁忙,请稍后重试");
}
try {
return orderservice.create(dto);
} finally {
ordersemaphore.release();
}
}
}
四、信号量 vs 锁 vs 线程池
| 工具 | 并发数 | 用途 | 区别 |
|---|---|---|---|
| lock/synchronized | 1 | 互斥访问 | 信号量(1) 的特例 |
| semaphore | n | 控制并发度 | 不关心是哪个线程释放 |
| 线程池 | n | 控制执行线程数 | 管理线程生命周期 |
关键区别:
// 锁:谁加的锁谁释放 lock.lock(); // ... 只能当前线程 unlock lock.unlock(); // 信号量:任何线程都可以释放 semaphore.acquire(); // 线程a 获取 // ... semaphore.release(); // 线程b 也可以释放(不要求同一线程)
这个特性使得信号量适合"生产者-消费者"场景:一个线程 acquire,另一个线程 release。
五、信号量的变体
1. 二元信号量(binary semaphore)
semaphore mutex = new semaphore(1); // 许可数=1,等效于互斥锁
与 lock 的区别:
- lock 有所有权(只能由持有者释放)
- 二元信号量无所有权(任何线程可释放)
2. 计数信号量(counting semaphore)
semaphore pool = new semaphore(10); // 标准用法
3. 带超时的信号量
// 超时未获取则放弃
boolean acquired = semaphore.tryacquire(3, timeunit.seconds);
if (!acquired) {
// 超时处理
}
4. 可增减的信号量
// 动态增加许可(如动态扩容连接池) semaphore.release(5); // 增加 5 个许可 // 动态减少许可 semaphore.acquire(3); // 消耗 3 个许可(不释放 = 永久减少)
六、分布式信号量
为什么需要分布式信号量
jdk semaphore 只在单个 jvm 内有效:
实例a: semaphore(10) → 允许 10 个 实例b: semaphore(10) → 允许 10 个 实际并发:可能 20 个同时访问(每个实例各 10 个)
分布式信号量通过 redis 等中间件共享计数器,所有实例共享同一个许可池。
redisson 分布式信号量
基本用法
@resource
private redissonclient redissonclient;
public void accessexternalapi() {
// 获取分布式信号量(所有实例共享)
rsemaphore semaphore = redissonclient.getsemaphore("semaphore:external-api");
// 首次需要设置许可数(只需执行一次)
semaphore.trysetpermits(10);
try {
// 获取许可(跨实例控制并发总数为 10)
semaphore.acquire();
callexternalapi();
} catch (interruptedexception e) {
thread.currentthread().interrupt();
} finally {
semaphore.release();
}
}
带超时的获取
rsemaphore semaphore = redissonclient.getsemaphore("semaphore:db-connection");
semaphore.trysetpermits(20);
// 最多等待 5 秒
boolean acquired = semaphore.tryacquire(5, timeunit.seconds);
if (acquired) {
try {
querydatabase();
} finally {
semaphore.release();
}
} else {
throw new businessexception("系统繁忙");
}
批量获取
// 一次获取 3 个许可(批量操作场景)
semaphore.acquire(3);
try {
batchprocess();
} finally {
semaphore.release(3);
}
redisson 过期信号量(permitexpirablesemaphore)
普通信号量的问题:如果获取许可后进程崩溃,许可永远不会被释放(许可泄漏)。
// 过期信号量:许可有 ttl,超时自动归还
rpermitexpirablesemaphore semaphore =
redissonclient.getpermitexpirablesemaphore("semaphore:task-runner");
semaphore.trysetpermits(5);
// 获取许可,10秒后自动释放(返回许可id)
string permitid = semaphore.acquire(10, timeunit.seconds);
try {
runtask();
// 手动提前释放
semaphore.release(permitid);
} catch (exception e) {
// 即使不释放,10秒后也会自动归还
semaphore.release(permitid);
}
对比:
| 类型 | 崩溃后许可 | 用法 |
|---|---|---|
| rsemaphore | 永久丢失(需人工恢复) | 稳定进程 |
| rpermitexpirablesemaphore | 超时自动归还 | 不可靠进程 |
七、分布式信号量的 redis 实现原理
数据结构
key: semaphore:external-api type: string value: 10(当前可用许可数)
acquire 操作(lua 脚本)
-- keys[1] = 信号量 key
-- argv[1] = 要获取的许可数
local permits = tonumber(redis.call('get', keys[1]))
if permits ~= nil and permits >= tonumber(argv[1]) then
-- 许可足够,扣减
redis.call('decrby', keys[1], argv[1])
return 1
end
-- 许可不足
return 0release 操作(lua 脚本)
-- 归还许可
redis.call('incrby', keys[1], argv[1])
-- 通知等待者
redis.call('publish', keys[2], argv[1])
return 1等待机制
获取失败时不轮询,使用 redis pub/sub 等待通知:
线程a acquire 失败
│
├─ 订阅 channel: redisson_sc:{semaphore:external-api}
│
├─ 阻塞等待通知
│
线程b release → publish 消息到 channel
│
└─ 线程a 收到通知 → 再次尝试 acquire八、实战场景
场景1:控制第三方 api 调用并发数
@service
public class thirdpartyapiservice {
@resource
private redissonclient redissonclient;
/**
* 第三方限制最多 5 个并发请求.
*/
public apiresponse callthirdpartyapi(apirequest request) {
rsemaphore semaphore = redissonclient.getsemaphore("semaphore:third-party-api");
semaphore.trysetpermits(5);
boolean acquired = false;
try {
acquired = semaphore.tryacquire(10, timeunit.seconds);
if (!acquired) {
throw new businessexception("第三方接口繁忙,请稍后重试");
}
return httpclient.post(request);
} catch (interruptedexception e) {
thread.currentthread().interrupt();
throw new businessexception("操作被中断");
} finally {
if (acquired) {
semaphore.release();
}
}
}
}
场景2:分布式限流(令牌桶简化版)
@component
public class distributedratelimiter {
@resource
private redissonclient redissonclient;
/**
* 每秒最多处理 100 个请求(所有实例合计).
*/
public boolean tryacquire(string resource) {
rsemaphore semaphore = redissonclient.getsemaphore("rate:" + resource);
return semaphore.tryacquire();
}
/**
* 每秒补充许可(定时任务).
*/
@scheduled(fixedrate = 1000)
public void refillpermits() {
rsemaphore semaphore = redissonclient.getsemaphore("rate:order-api");
int current = semaphore.availablepermits();
if (current < 100) {
semaphore.release(100 - current); // 补充到 100
}
}
}
场景3:数据库连接池保护
@service
public class databaseservice {
@resource
private redissonclient redissonclient;
// 数据库最大连接 50,预留 10 给管理操作
// 业务最多使用 40 个连接
private static final int max_biz_connections = 40;
public <t> t executequery(supplier<t> query) {
rsemaphore semaphore = redissonclient.getsemaphore("semaphore:db-biz-conn");
semaphore.trysetpermits(max_biz_connections);
try {
semaphore.acquire();
return query.get();
} catch (interruptedexception e) {
thread.currentthread().interrupt();
throw new runtimeexception(e);
} finally {
semaphore.release();
}
}
}
场景4:并行任务控制
/**
* 导出报表:允许系统同时最多处理 3 个导出任务(防止 oom).
*/
@service
public class reportexportservice {
@resource
private redissonclient redissonclient;
public void exportreport(integer reportid) {
rpermitexpirablesemaphore semaphore =
redissonclient.getpermitexpirablesemaphore("semaphore:report-export");
semaphore.trysetpermits(3);
string permitid = null;
try {
// 获取许可,最多持有 5 分钟(防止任务卡死占用许可)
permitid = semaphore.tryacquire(30, 300, timeunit.seconds);
if (permitid == null) {
throw new businessexception("当前导出任务过多,请稍后再试");
}
doexport(reportid);
} finally {
if (permitid != null) {
semaphore.release(permitid);
}
}
}
}
九、信号量 vs 其他并发控制工具
| 工具 | 控制维度 | 适用场景 |
|---|---|---|
| semaphore | 并发数量 | 控制同时执行的操作数 |
| ratelimiter | 速率(每秒n个) | 控制请求频率 |
| lock | 互斥(0或1) | 独占资源 |
| countdownlatch | 等待计数归零 | 等待多个任务完成 |
| cyclicbarrier | 等待n个线程到达 | 多线程同步汇合 |
| 线程池 | 工作线程数 | 管理执行资源 |
信号量 vs 线程池
// 线程池方式:控制执行线程数
executorservice pool = executors.newfixedthreadpool(10);
pool.submit(() -> callapi()); // 超过 10 个则排队
// 信号量方式:控制并发数(不管你用什么线程)
semaphore sem = new semaphore(10);
sem.acquire();
try {
callapi(); // 可以在任何线程中执行
} finally {
sem.release();
}
区别:
- 线程池管理线程的创建和销毁
- 信号量只管"允许多少个同时执行",不管线程来自哪里
两者常配合使用:线程池控制线程总量,信号量控制某类操作的并发量。
信号量 vs ratelimiter
// 信号量:同一时刻最多 10 个并发 // 如果每个请求处理 1 秒,则吞吐约 10/s semaphore sem = new semaphore(10); // ratelimiter:每秒最多 10 个请求 // 不管并发多少,严格控制速率 ratelimiter limiter = ratelimiter.create(10.0); limiter.acquire(); // 平滑限流
| 信号量 | ratelimiter | |
|---|---|---|
| 控制的是 | 同时进行的数量 | 单位时间的数量 |
| 短时间突发 | 允许(并发数内) | 平滑(不允许突发) |
| 请求处理时间影响 | 处理越慢,吞吐越低 | 不受处理时间影响 |
十、注意事项与陷阱
陷阱1:许可泄漏
// 错误:异常时不释放许可
semaphore.acquire();
dowork(); // 如果这里抛异常
semaphore.release(); // 这行不会执行 → 许可永久丢失
// 正确:finally 中释放
semaphore.acquire();
try {
dowork();
} finally {
semaphore.release();
}
陷阱2:释放多于获取
semaphore sem = new semaphore(3); // 没有 acquire 就 release → 许可数变成 4! sem.release(); // 现在许可数 = 4,超过了设计的 3
信号量不会校验"是否之前获取过",多余的 release 会增加许可总数。
陷阱3:分布式环境下的初始化竞争
// 多个实例同时启动,都执行 trysetpermits semaphore.trysetpermits(10); // 实例a semaphore.trysetpermits(10); // 实例b(如果 key 已存在则不生效)
trysetpermits 是"不存在才设置"(类似 setnx),所以多实例并发调用是安全的。但如果要修改许可数,需要用 addpermits:
// 扩容:增加 5 个许可 semaphore.addpermits(5); // 缩容:减少 3 个许可(当前可用 >= 3 才能成功) semaphore.addpermits(-3);
陷阱4:try-finally 中的 acquire 返回值
// 错误:tryacquire 返回 false 但 finally 仍然 release
boolean acquired = semaphore.tryacquire();
try {
if (acquired) dowork();
} finally {
semaphore.release(); // acquired=false 时多释放了!
}
// 正确:条件释放
boolean acquired = semaphore.tryacquire();
try {
if (!acquired) {
throw new businessexception("繁忙");
}
dowork();
} finally {
if (acquired) {
semaphore.release();
}
}
十一、总结
| 概念 | 一句话 |
|---|---|
| 信号量 | 一个计数器,控制"同时有多少个"能执行 |
| acquire | 获取许可(计数-1),许可为0时阻塞 |
| release | 释放许可(计数+1),唤醒等待者 |
| 与锁的区别 | 锁是二元的(0或1),信号量是n元的 |
| 分布式信号量 | 用 redis 存储计数,所有实例共享许可池 |
| 过期信号量 | 许可有 ttl,进程崩溃后自动归还 |
| 核心价值 | 保护有限资源不被过度并发访问 |
| 常见用途 | api 并发控制、连接池保护、任务并行度限制 |
到此这篇关于java中信号量semaphore的使用的文章就介绍到这了,更多相关java 信号量semaphore内容请搜索代码网以前的文章或继续浏览下面的相关文章希望大家以后多多支持代码网!
发表评论