信号量模式(Semaphore Pattern)
一、核心思想
信号量是一种并发控制机制,用于限制同时访问某个资源的 goroutine 数量。它是 Mutex(只允许 1 个)和无限制(全部允许)之间的中间地带——允许 N 个 goroutine 并行执行。
Go 中实现信号量有两种方式:
- Buffered Channel:
chan struct{}作为令牌桶,容量 = 最大并发数 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(¤tRunning, 1)fmt.Printf("任务#%d 开始执行 (当前并发: %d)\n", taskID, running)time.Sleep(500 * time.Millisecond) // 模拟工作atomic.AddInt32(¤tRunning, -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)
// 执行工作...
七、实践要点总结
- acquire 放在
go之前:限制 goroutine 数量本身,而不仅仅是执行数 - defer release:panic 也能保证归还令牌
- 不要 close 信号量 channel:关闭后继续 send 会 panic,没必要关闭
- 权重配对:
Acquire(ctx, w)和Release(w)的权重必须一致,否则配额计算会错乱 - acquire 权重超过总容量会永久阻塞:除非 context 取消
八、学习小结
信号量是 Go 并发编程中最实用的模式之一。核心知识:
- Buffered Channel 是最简单的方式:
chan struct{}容量 = 并发上限 semaphore.Weighted提供权重 + context 支持,适合复杂场景- acquire 必须在
go之前,否则 goroutine 本身不会被限制 - 信号量限并发,不限速率——速率用
rate.Limiter