Appearance
34|Worker Pool 模式
批量巡检 1000 台机器时,如果同时发起 1000 个请求,可能把本机端口、DNS、目标服务或中间网关打满。worker pool 用固定数量的 worker 控制并发上限,避免所有目标同时发起请求。
一、基础结构
go
package main
import (
"fmt"
"sync"
"time"
)
type Result struct {
Host string
OK bool
}
func check(host string) Result {
time.Sleep(300 * time.Millisecond)
return Result{Host: host, OK: true}
}
func worker(jobs <-chan string, results chan<- Result, wg *sync.WaitGroup) {
defer wg.Done()
for host := range jobs {
results <- check(host)
}
}
func main() {
hosts := []string{"web01", "web02", "web03", "db01", "db02"}
concurrency := 2
jobs := make(chan string)
results := make(chan Result)
var wg sync.WaitGroup
for i := 0; i < concurrency; i++ {
wg.Add(1)
go worker(jobs, results, &wg)
}
go func() {
for _, host := range hosts {
jobs <- host
}
close(jobs)
}()
go func() {
wg.Wait()
close(results)
}()
for result := range results {
fmt.Println(result.Host, result.OK)
}
}这个结构里 jobs 是任务队列,results 是结果队列,worker 数量就是并发上限。并发数通常根据目标服务承受能力、网络环境和超时时间调整。内网 HTTP 健康检查可以稍高,SSH 登录、数据库探测这类操作通常更保守。
二、执行顺序
- 主流程启动固定数量 worker
- 投递协程把所有 host 写入
jobs - worker 从
jobs取任务,执行check,把结果写入results jobs关闭后,worker 的for host := range jobs结束- 所有 worker 结束后关闭
results - 主流程读取完
results后退出
jobs 和 results 的关闭顺序很关键。jobs 由投递方关闭,表示没有新任务;results 由等待 worker 结束的协程关闭,表示没有新结果。提前关闭 results,worker 再写结果会 panic。
三、带 context 的 worker pool
实际生产代码需要支持超时和取消:
go
func worker(ctx context.Context, jobs <-chan string, results chan<- Result, wg *sync.WaitGroup) {
defer wg.Done()
for {
select {
case host, ok := <-jobs:
if !ok {
return
}
results <- check(ctx, host)
case <-ctx.Done():
return
}
}
}worker 同时监听 jobs 和 ctx.Done(),收到取消信号时立即退出,不再执行新任务。
四、动态 worker 数量
固定 worker 数量不够灵活。可以按需伸缩:任务堆积时增加 worker,空闲时减少:
go
type Pool struct {
jobs chan string
results chan Result
workers int
mu sync.Mutex
}
func (p *Pool) Scale(target int) {
p.mu.Lock()
defer p.mu.Unlock()
for p.workers < target {
p.workers++
go worker(p.jobs, p.results)
}
}动态伸缩需要仔细处理 worker 的生命周期,避免资源泄漏。大多数场景下固定数量更稳定。
五、常见错误
忘记关闭 jobs
投递方不关闭 jobs,worker 的 range 永远不会结束,wg.Wait() 永远等不到,程序卡住。
用 Mutex 保护 channel
channel 本身是并发安全的,不需要加锁。给 channel 操作加 Mutex 是多余的,还会降低性能。
worker 里 panic 没有被恢复
go
func worker(jobs <-chan string, results chan<- Result, wg *sync.WaitGroup) {
defer wg.Done()
for host := range jobs {
results <- check(host) // 如果 check panic,wg.Done() 不会执行
}
}在 worker 开头加 recover 保护,或者外层用 sync.WaitGroup 时设置超时兜底。
结果不消费导致死锁
go
for _, host := range hosts {
jobs <- host
}如果 jobs 是无缓冲 channel,且所有 worker 都在等 results 被消费,而主流程还没开始读 results,就会形成循环等待。确保有独立的 goroutine 消费结果,或给 results 足够大的缓冲。