Skip to content

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 登录、数据库探测这类操作通常更保守。

二、执行顺序

  1. 主流程启动固定数量 worker
  2. 投递协程把所有 host 写入 jobs
  3. worker 从 jobs 取任务,执行 check,把结果写入 results
  4. jobs 关闭后,worker 的 for host := range jobs 结束
  5. 所有 worker 结束后关闭 results
  6. 主流程读取完 results 后退出

jobsresults 的关闭顺序很关键。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 同时监听 jobsctx.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 足够大的缓冲。