Go worker pool
用固定数量 worker 消费任务 channel 控制并发度;所有权、收尾、取消与背压,以及与 pipeline 的分工边界。
[!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。
关闭纪律:
- 生产者关
jobs(没有新图了) - 等所有 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(顺序可能打乱)。
结合场景再看三个关注点
- worker 数 = 资源上限,不是任务数;任务再多也只开 3 个工人。
close(jobs)只有生产者做;工人只消费。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-context 与 go-goroutine-channel-context-and-cancellation-relationship。
错误与部分失败
常见策略:
- errgroup:任一失败 cancel 全组,适合“全部成功才有意义”。
- 结果结构体带 error:适合“尽量完成,汇总失败项”。
- 独立 error channel:要严格所有权与关闭顺序,复杂度更高。
动态池?
Go 标准库没有 Java 式 ThreadPoolExecutor。多数服务用:
- 固定 N(=CPU 或下游限额)
- 或信号量(
chan struct{}/semaphore)限制在途 goroutine
“动态扩缩”通常是运维层(副本数)问题,而不是进程内无限加 worker。
边界情况
- jobs 不关闭:worker
range jobs永久阻塞 → 泄漏。 - results 被 worker 各自 close:panic 或丢数据。
- main 不等待 worker:进程退出截断未完成任务。
- 缓冲过大:伪装成无界队列,延迟暴露背压。
- 任务内再无限
go:池外泄漏,上限形同虚设。 - 共享可变状态:池只限并发数,不自动消除 data race(需
-race与同步)。 - 取消后仍向 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 做成全局可变单例乱投递 错误:无超时、无背压、无关闭协议的全局队列。
正确:生命周期绑定服务/请求范围,协议写进类型与文档。
工程实践
- 签名表达所有权:
jobs <-chan Job、results chan<- Result。 - 用 context 贯穿处理函数,
defer cancel()在创建点。 - 优雅停机:停止接收新请求 → cancel/关闭 jobs → Wait → 再关 listener(go-graceful-shutdown)。
- 指标:队列长度、在途数、处理时延、失败率(go-observability)。
- 任务幂等:至少一次投递时下游要可重试安全。
- 测试:用小 N、可控 jobs、fake clock 或即时完成任务;加
-race。 - 优先标准组合:
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。
自测题
概念题
- worker pool 主要限制的是什么资源行为?
- 谁应该
close(jobs)?谁应该close(results)? - 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 数、队列缓冲,以及客户端断开时如何停?
参考答案
展开
- 同时执行的任务数(及由此绑定的 CPU/连接/FD 等)。
- jobs 由任务发送方/拥有者关闭;results 由“确认所有发送 results 的 goroutine 已结束”的协调者关闭。
- 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 博文术语)