当前位置: 首页 > news >正文

信号量模式(Semaphore Pattern)

信号量模式(Semaphore Pattern)

一、核心思想

信号量是一种并发控制机制,用于限制同时访问某个资源的 goroutine 数量。它是 Mutex(只允许 1 个)和无限制(全部允许)之间的中间地带——允许 N 个 goroutine 并行执行。

Go 中实现信号量有两种方式:

  1. Buffered Channelchan struct{} 作为令牌桶,容量 = 最大并发数
  2. golang.org/x/sync/semaphore:官方扩展库,支持权重和 context

二、为什么需要信号量

无限制并发的风险

// ❌ 危险:500 个 URL 同时请求
for _, url := range urls {go fetch(url) // 500 个 goroutine 同时打开 socket
}

后果:

  • 文件描述符耗尽too many open files
  • 下游限流封禁:第三方 API 返回 429
  • 内存暴涨:500 个 goroutine 各自持有响应缓冲区
  • 级联超时:连接池打满 → 排队 → 超时 → 重试 → 更糟

goroutine 很轻(初始 2KB 栈),但它们触碰的资源不轻

三、Buffered Channel 信号量

基本原理

容量为 N 的 buffered channel = 容许 N 个并发sem <- struct{}{}  // 获取令牌(池满则阻塞)
<-sem              // 归还令牌

核心代码模式

sem := make(chan struct{}, maxConcurrent)for _, item := range items {sem <- struct{}{}  // ① 获取令牌(在 goroutine 外面!)go func(it Item) {defer func() { <-sem }()  // ② defer 归还令牌process(it)}(item)
}

关键陷阱:acquire 放在哪里

// ✅ 正确:acquire 在 go 语句之前
for _, item := range items {sem <- struct{}{}  // 主 goroutine 阻塞,控制 goroutine 数量go func() {defer func() { <-sem }()process(item)}()
}// ❌ 错误:acquire 在 goroutine 内部
for _, item := range items {go func() {sem <- struct{}{}  // 所有 goroutine 都已启动,只是排队等令牌defer func() { <-sem }()process(item)}()
}

区别:放在 go 之前,限制的是活跃 goroutine 数量;放在 go 之后,限制的是并发执行数量但 goroutine 已经全部创建了(500 个 goroutine 在等令牌,内存照样涨)。

四、Go 实现示例

方式一:Buffered Channel 信号量(限制 HTTP 并发)

package mainimport ("fmt""sync""sync/atomic""time"
)func main() {// 模拟 20 个任务,限制最多 4 个同时执行const maxConcurrent = 4const totalTasks = 20sem := make(chan struct{}, maxConcurrent)var wg sync.WaitGroup// 用原子计数器追踪当前并发数var currentRunning int32for i := 1; i <= totalTasks; i++ {wg.Add(1)sem <- struct{}{} // 获取令牌,满 4 个后阻塞go func(taskID int) {defer wg.Done()defer func() { <-sem }() // 归还令牌running := atomic.AddInt32(&currentRunning, 1)fmt.Printf("任务#%d 开始执行 (当前并发: %d)\n", taskID, running)time.Sleep(500 * time.Millisecond) // 模拟工作atomic.AddInt32(&currentRunning, -1)}(i)}wg.Wait()fmt.Println("\n所有任务完成!")
}

方式二:加权信号量(semaphore.Weighted

适用于不同任务消耗不同资源配额的场景。例如:小任务占 1 个配额,大任务占 3 个配额。

package mainimport ("context""fmt""sync""sync/atomic""time""golang.org/x/sync/semaphore"
)func main() {// 总容量 10 个权重单位// 大任务占 5,中任务占 3,小任务占 1sem := semaphore.NewWeighted(10)var wg sync.WaitGroupvar running int32tasks := []struct {name   stringweight int64}{{"小任务A", 1},{"小任务B", 1},{"大任务", 5},{"中任务", 3},{"小任务C", 1},{"大任务2", 5},{"小任务D", 1},{"中任务2", 3},}for _, task := range tasks {wg.Add(1)// 获取配额,支持 context 取消if err := sem.Acquire(context.Background(), task.weight); err != nil {fmt.Printf("%s: 获取配额失败: %v\n", task.name, err)wg.Done()continue}go func(name string, weight int64) {defer wg.Done()defer sem.Release(weight)r := atomic.AddInt32(&running, 1)fmt.Printf("[开始] %-8s 权重=%d 当前并发=%d\n", name, weight, r)time.Sleep(800 * time.Millisecond) // 模拟工作atomic.AddInt32(&running, -1)fmt.Printf("[完成] %-8s 权重=%d\n", name, weight)}(task.name, task.weight)}wg.Wait()// 用获取全部容量的方式等待所有工作完成// 这个技巧可以代替 WaitGroupfmt.Println("\n所有任务完成!")
}

方式三:带 context 取消的信号量

package mainimport ("context""fmt""time"
)func main() {// 模拟 3 秒超时限制下的并发控制// 即使有 100 个任务,只允许 2 个同时跑,3 秒后取消所有等待sem := make(chan struct{}, 2)ctx, cancel := context.WithTimeout(context.Background(), 3*time.Second)defer cancel()done := make(chan struct{})// 启动 5 个任务for i := 1; i <= 5; i++ {go func(id int) {select {case sem <- struct{}{}:defer func() { <-sem }()fmt.Printf("任务#%d 获得令牌\n", id)time.Sleep(1 * time.Second)fmt.Printf("任务#%d 完成\n", id)case <-ctx.Done():fmt.Printf("任务#%d 超时取消: %v\n", id, ctx.Err())}done <- struct{}{}}(i)}// 等待所有任务结束for i := 0; i < 5; i++ {<-done}fmt.Println("全部结束")
}

五、信号量 vs Worker Pool vs 限流器

机制 限制维度 典型场景 Go 实现
信号量 最大并发数 限制同时打开的连接数 chan struct{} / semaphore.Weighted
Worker Pool 固定 worker 数量 长期运行的任务队列 N 个 goroutine 读同一个 channel
限流器 时间窗口内的频率 API 调用频率控制 golang.org/x/time/rate

信号量 ≠ 限流器:信号量限制"同时有多少个在跑",限流器限制"每秒能跑多少个"。

可以组合使用:

// 同时限制并发数和速率
sem := make(chan struct{}, 3)              // 最多 3 个并发
limiter := rate.NewLimiter(rate.Limit(10), 1) // 每秒最多 10 次for _, job := range jobs {limiter.Wait(ctx)   // 速率门控sem <- struct{}{}   // 并发门控go func() {defer func() { <-sem }()process(job)}()
}

六、semaphore.Weighted 的 TryAcquire

非阻塞获取,适合过载丢弃而不是排队的场景:

if !sem.TryAcquire(1) {// 配额已满,直接拒绝或降级处理return errors.New("系统繁忙,请稍后重试")
}
defer sem.Release(1)
// 执行工作...

七、实践要点总结

  1. acquire 放在 go 之前:限制 goroutine 数量本身,而不仅仅是执行数
  2. defer release:panic 也能保证归还令牌
  3. 不要 close 信号量 channel:关闭后继续 send 会 panic,没必要关闭
  4. 权重配对Acquire(ctx, w)Release(w) 的权重必须一致,否则配额计算会错乱
  5. acquire 权重超过总容量会永久阻塞:除非 context 取消

八、学习小结

信号量是 Go 并发编程中最实用的模式之一。核心知识:

  • Buffered Channel 是最简单的方式:chan struct{} 容量 = 并发上限
  • semaphore.Weighted 提供权重 + context 支持,适合复杂场景
  • acquire 必须在 go 之前,否则 goroutine 本身不会被限制
  • 信号量限并发,不限速率——速率用 rate.Limiter
http://www.jsqmd.com/news/1371543/

相关文章:

  • 2026南京专业漏水检测公司推荐本地正规堵漏公司服务电话 - 知途管道科技
  • 源于国赛标准,归于日常课堂——哈弗M6发动机仿真教学软件
  • 债权转让公告该怎么登报?报纸要满足哪些条件?报纸选用标准整理 - 点办通
  • Unity字体字符集全解析:解决中文乱码与7000汉字支持方案
  • 抗变形烤盘哪家供应商好
  • 手持光谱仪购玉避坑指南:玉石从业者与玩家必备检测工具 - 巴斯德仪器
  • Simulink建模:双馈风机在英格兰10机39节点系统的应用
  • 多进程并行化BPE分词器实现:从算法原理到工程优化
  • 【单片机毕设案例分享】基于 STM32 的温光双传感联动智能窗帘控制器开发 基于 STM32 单片机的参数可配置室内智能窗帘硬件系统实现(018202)
  • C#、C++、Java三种语言实现塔防游戏:架构、性能与跨语言开发对比
  • Android Studio 2021.1.1 稳定版安装与配置全指南:从环境搭建到高效开发
  • 免费调用Kimi与GLM-5.2 API:开源代理部署与实战指南
  • 2026 安捷伦2ml样品瓶批发拿货渠道,现货充足专业代理商福州希音生物 - 品牌推荐大师
  • 2026年8月geo优化公司全景盘点:实测数据盘点与分级选型建议 - 天下观知
  • 2026 年唯美汽车服务|喀什车灯升级改装选购避坑完整干货 - 收录优先
  • WordPress媒体文件自动清理插件开发指南
  • 大件家具退货率从35%降到20%:我把英语客服外包后的真实体验
  • Ubuntu根目录磁盘空间不足?LVM在线扩容实战指南
  • 虚幻引擎VDB插件实战:从原理到性能优化的体积渲染指南
  • 从零构建低代码平台:可视化应用构建原型实践指南
  • 基于OpenClaw与SecGPT-14B的智能开源组件审计与SBOM自动生成实践
  • 磁吸悬浮键盘深度评测:iPad高效输入与手写批注一体化解决方案
  • 解决Python包安装中的Metadata-Version冲突问题
  • Linux磁盘IO性能监控与优化:从iostat到实战场景解析
  • AI API额度监控与自动化调度系统实战:告别额度耗尽困扰
  • 终极Sketch设计标注革命:如何用MeaXure实现设计开发无缝协作
  • 债权转让公告登报:债权转让公告登报怎么做才算具备合规效力? - 点办通
  • WorkBuddy入门到精通:零基础小白也能成为AI效率高手
  • 学术生产力提升的实用路径与高效方法探索
  • 前后端分离医院资源管理系统系统|SpringBoot+Vue+MyBatis+MySQL完整源码+部署教程