ARTICLE DETAIL

建站实战干货

来自一线的建站与推广经验沉淀,每一条都经过真实交付验证。

信号量模式(Semaphore Pattern)

2026/8/11 9:17:45 拓冰建站 浏览量
信号量模式(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