Go worker pool

用固定数量 worker 消费任务 channel 控制并发度;所有权、收尾、取消与背压,以及与 pipeline 的分工边界。

#type / concept #status / growing #tech / dev #resource / go

[!info] 关联笔记

Go worker pool

这个概念为什么会出现

go process(job) 写起来极便宜,但工程上会立刻撞到:

  • 任务量远大于 CPU / 连接池 / 下游 QPS 能承受的并发
  • 瞬时尖峰把进程打成成千上万 goroutine,调度与内存抖动
  • 需要明确上限、可观测的在途任务数、可取消的收尾

worker pool 的答案是:预先(或按需但有上限地)启动固定数量的 worker,从共享任务队列取活。它不发明新语言特性,而是把“并发度控制 + 任务所有权 + 收尾”组合成稳定模式。官方在 Go Concurrency Patterns: Pipelines 里用 fan-out / fan-in 展示了同族思想:多 goroutine 从同一输入 channel 取任务,再汇总结果。

[!abstract] 一句话理解 worker pool = 有界 worker 集合 + 任务输入 channel(可选结果 channel)+ 明确的关闭/取消/等待协议;目标是控制并发资源,而不是“开更多 goroutine 就更快”。

最小可运行示例

先把示例放进业务场景,再看代码:

场景:批量缩略图任务,固定 3 个 worker

运营上传 5 张图,每张要做一次 CPU/IO 混合处理(这里用 j*j 假装“算缩略图”)。
若每张图 go process(),高峰会把机器打满。

worker pool:固定 3 个工人抢 jobs 队列;结果写入 results
关闭纪律:

  1. 生产者关 jobs(没有新图了)
  2. 等所有 worker 退出后再关 results(不能让某个 worker 提前 close)
package main

import (
	"fmt"
	"sync"
)

func main() {
	// jobs:待处理图片 ID 队列(缓冲 8 = 有限排队,满则背压)
	jobs := make(chan int, 8)
	// results:处理结果(演示用平方值)
	results := make(chan int, 8)

	const workers = 3 // 并发上限:同时最多 3 张在处理
	var wg sync.WaitGroup
	for w := 0; w < workers; w++ {
		wg.Add(1)
		go func(workerID int) {
			defer wg.Done()
			// range jobs:队列关闭且排空后退出
			for jobID := range jobs {
				// 业务:生成缩略图;这里用平方假装 CPU 工作
				_ = workerID
				results <- jobID * jobID
			}
		}(w)
	}

	// 协调者:投递任务 → 关 jobs → 等工人 → 关 results
	go func() {
		for j := 1; j <= 5; j++ {
			jobs <- j // 5 张图
		}
		close(jobs) // 发送方关闭:工人的 range 才能结束
		wg.Wait()   // 所有工人不再写 results
		close(results)
	}()

	// 主流程:收齐全部结果(顺序不保证)
	for r := range results {
		fmt.Println("thumb result:", r)
	}
}

建议运行:

go run .

期望:打印 5 个结果,集合为 1 4 9 16 25(顺序可能打乱)。

结合场景再看三个关注点

  1. worker 数 = 资源上限,不是任务数;任务再多也只开 3 个工人。
  2. close(jobs) 只有生产者做;工人只消费。
  3. close(results) 在所有发送者结束后由协调者做,避免“close of closed channel”或丢结果。

核心概念与准确模型

组成部件

部件职责
job channel任务队列;缓冲大小决定排队深度与背压
worker 集合固定 N 个循环:取任务 → 处理 → 可选写结果
result channel / 回调输出路径;所有权必须单一
WaitGroup / errgroup等待 worker 退出
context取消进行中的工作、停止再取新任务

与 fan-out 的关系

在 pipelines 博文中:

  • fan-out:多个函数从同一 channel 读(竞争消费)
  • fan-in:多个 channel 汇入一个

经典 worker pool 就是 fan-out 消费 jobs +(可选)fan-in 汇总 results。不必把“pool 对象”做成框架,channel 本身就是队列。

有界 vs 无界

无界:每个任务 go f(job)     → 并发度 = 任务数(危险)
有界:N workers + jobs 队列  → 并发度 ≤ N

缓冲 jobs 的容量是第二道闸门:满了则生产者阻塞(背压),而不是静默堆积到 OOM。

取消语义

func worker(ctx context.Context, jobs <-chan Job, out chan<- Result) {
	for {
		select {
		case <-ctx.Done():
			return
		case j, ok := <-jobs:
			if !ok {
				return
			}
			// 处理时也要能被取消:优先用支持 ctx 的 API
			r, err := do(ctx, j)
			if err != nil {
				// 约定:如何上报错误(result 内嵌 / 独立 err ch / errgroup)
				continue
			}
			select {
			case out <- r:
			case <-ctx.Done():
				return
			}
		}
	}
}

取消不是“杀死 goroutine”,而是协作退出:停取新任务、中断可取消的 I/O、尽快释放锁与连接。详见 go-contextgo-goroutine-channel-context-and-cancellation-relationship

错误与部分失败

常见策略:

  1. errgroup:任一失败 cancel 全组,适合“全部成功才有意义”。
  2. 结果结构体带 error:适合“尽量完成,汇总失败项”。
  3. 独立 error channel:要严格所有权与关闭顺序,复杂度更高。

动态池?

Go 标准库没有 Java 式 ThreadPoolExecutor。多数服务用:

  • 固定 N(=CPU 或下游限额)
  • 或信号量(chan struct{} / semaphore)限制在途 goroutine

“动态扩缩”通常是运维层(副本数)问题,而不是进程内无限加 worker。

边界情况

  1. jobs 不关闭:worker range jobs 永久阻塞 → 泄漏。
  2. results 被 worker 各自 close:panic 或丢数据。
  3. main 不等待 worker:进程退出截断未完成任务。
  4. 缓冲过大:伪装成无界队列,延迟暴露背压。
  5. 任务内再无限 go:池外泄漏,上限形同虚设。
  6. 共享可变状态:池只限并发数,不自动消除 data race(需 -race 与同步)。
  7. 取消后仍向 out 发送:接收方已退出 → 发送阻塞 → 二次泄漏。

常见误区

[!warning] 常见误区:worker pool 是默认并发方案 错误:所有并发都套 pool。
正确:独立短任务且自然有界可用 errgroup;阶段流水线用 pipeline;共享结构用 mutex。

[!warning] 常见误区:N 越大越快 错误:workers = 10000 处理 CPU 密集任务。
正确:CPU 密集常近 GOMAXPROCS;I/O 密集按下游配额与延迟实验定。

[!warning] 常见误区:关闭 = 取消 错误:只 close jobs 就以为进行中的 HTTP 会停。
正确:close 停“新任务”;进行中的工作要靠 context / 可中断 API。

[!warning] 常见误区:把 pool 做成全局可变单例乱投递 错误:无超时、无背压、无关闭协议的全局队列。
正确:生命周期绑定服务/请求范围,协议写进类型与文档。

工程实践

  1. 签名表达所有权jobs <-chan Jobresults chan<- Result
  2. 用 context 贯穿处理函数,defer cancel() 在创建点。
  3. 优雅停机:停止接收新请求 → cancel/关闭 jobs → Wait → 再关 listener(go-graceful-shutdown)。
  4. 指标:队列长度、在途数、处理时延、失败率(go-observability)。
  5. 任务幂等:至少一次投递时下游要可重试安全。
  6. 测试:用小 N、可控 jobs、fake clock 或即时完成任务;加 -race
  7. 优先标准组合errgroup.Group + 信号量 往往比自研“万能 Pool 结构体”更清晰。
// 信号量式有界并发(不必常驻 worker)
sem := make(chan struct{}, 8)
var wg sync.WaitGroup
for _, job := range all {
	wg.Add(1)
	sem <- struct{}{}
	go func(j Job) {
		defer wg.Done()
		defer func() { <-sem }()
		_ = handle(ctx, j)
	}(job)
}
wg.Wait()

常驻 pool 适合持续到达的任务流;信号量适合一批已知任务的扇出。

可验证实验

实验 1:关闭顺序

去掉 close(jobs),观察程序是否卡住。

实验 2:错误关闭 results

让每个 worker 处理完后 close(results),观察 panic。

实验 3:背压

jobs 无缓冲、worker 处理很慢,看生产者是否在发送处阻塞。

实验 4:取消

WithCancel 后 cancel,确认 worker 从 select 退出且不再写 results。

实验 5:对比无界 go

对 1e5 任务分别“每任务一个 goroutine”与“8 worker”,比较峰值内存与完成时间(本机观察即可)。

本节总结

  • 本质:有界并发的任务消费模式,不是魔法调度器。
  • 关键规则:发送方关闭输入;协调者关闭输出;取消与关闭分工;等待所有 worker。
  • 最易错:关闭权混乱、无取消的“假收尾”、无限缓冲伪装。
  • 下一步:阶段化数据流见 go-pipeline-pattern;取消总图见 go-goroutine-channel-context-and-cancellation-relationship

自测题

概念题

  1. worker pool 主要限制的是什么资源行为?
  2. 谁应该 close(jobs)?谁应该 close(results)
  3. close jobs 与 cancel context 分别停止什么?

代码推理题

jobs := make(chan int)
go func() {
	jobs <- 1
	// 忘记 close(jobs)
}()
for j := range jobs {
	fmt.Println(j)
}

程序能否退出?为什么?

工程思考题

HTTP 服务要把上传文件扇出到 50 个第三方 API。如何选 worker 数、队列缓冲,以及客户端断开时如何停?

参考答案

展开
  1. 同时执行的任务数(及由此绑定的 CPU/连接/FD 等)。
  2. jobs 由任务发送方/拥有者关闭;results 由“确认所有发送 results 的 goroutine 已结束”的协调者关闭。
  3. close jobs:不再派发新任务;cancel:进行中的可取消工作应停止,且循环应退出。
    代码题:不能正常结束,range jobs 在未关闭且无更多发送时阻塞(发送方已退出则死锁/泄漏)。
    工程题:worker 上限对齐第三方配额与本机连接池;小缓冲制造背压;请求 ctx 取消时停止读队列并取消 in-flight HTTP;用 errgroup 或结果汇总错误。

延伸阅读与资料来源

资料类型支撑
Go Concurrency Patterns: Pipelines官方博客fan-out/fan-in、取消、阶段组合
Go Concurrency Patterns: Context官方博客取消在调用树中的传播
Share Memory By Communicating官方博客channel 表达并发结构
errgroup扩展库错误与取消组
Effective Go — Concurrency文档并发风格

笔记元信息

  • 建议文件名:go-worker-pool.md
  • 所属阶段:并发模式
  • 建议下一篇:Go pipeline 模式
  • 本篇状态:已深化(结构完整;对齐 pipelines 博文术语)
创建于 2026/6/20 更新于 2026/7/15