10.3 worker pool 与 pipeline
10.2 节的任务队列只有一个消费者。这解决了「扫描与处理解耦」的问题,但没解决「处理太慢」的问题:如果每个到期任务都要写一次数据库,串行消费的速度会直接成为瓶颈。
最直觉的补救是「每个任务起一个 goroutine」。这确实快,但它把并发量交给了任务数量——某天积压了一万个任务,你就同时打开了一万个数据库连接。并发需要上限,worker pool 就是最常用的那个上限。
本节把 TaskAPI 推进到:到期任务队列由固定 4 个 worker 消费,批量关闭任务,并用 pipeline 把「扫描 → 排队 → 关闭 → 汇总」串成一条流水线。
10.3.1 为什么需要 worker pool
先看两种写法的实测差距。同一批 12 个任务,每个耗时 30 毫秒(Apple M1 Pro,go1.27.0):
| 写法 | 实测耗时 | 峰值并发 |
|---|---|---|
| 每个任务一个 goroutine | 31 ms | 12 |
| 4 个 worker 消费队列 | 93 ms | 4 |
无限制并发:12 个任务耗时 31ms
4 worker:共处理 12 个,耗时 93ms(单任务 30ms)
无限制并发快了 3 倍,代价是峰值并发不可控。在你的笔记本上 12 个 goroutine 毫无压力,但在生产环境里这个数字等于「积压任务数」,可能是十万。而每个 goroutine 背后都可能占着一个数据库连接、一个文件句柄、一次 HTTP 请求。
所以真正的取舍不是「快 vs 慢」,而是上游速度、下游承受力、资源可控性三者的平衡。结论:凡是任务会触碰外部资源(IO、网络、锁、连接池),就应该用 worker pool 限并发;纯 CPU 计算的任务交给 GOMAXPROCS 和运行时调度即可,不需要手工建池。
10.3.2 worker pool 的骨架
worker pool 的结构只有三部分:
jobs channel ← 生产者往里投递任务
│
├── worker 1 ──┐
├── worker 2 │ 各自 for range jobs,取一个处理一个
├── worker 3 │
└── worker 4 ──┘
results channel → 汇总结果
写成代码:
func worker(id int, jobs <-chan Task, results chan<- string, wg *sync.WaitGroup) {
defer wg.Done()
for t := range jobs {
time.Sleep(20 * time.Millisecond) // 模拟写库
results <- fmt.Sprintf("worker-%d 关闭任务 #%d(%s)", id, t.ID, t.Title)
}
}
三个要点:
- worker 的退出条件是
jobs被关闭。for range jobs会在 channel 关闭且缓冲读空后自动结束,不需要额外的退出信号。 - worker 自己不去关
jobs。它只是接收方,关闭是生产者的事(10.2 节的铁律)。 wg.Done()放在defer,保证无论 worker 从哪条路径退出都被计数。
主流程负责启动 worker、投递任务、关闭 jobs:
var wg sync.WaitGroup
const workers = 3
for i := 1; i <= workers; i++ {
wg.Add(1)
go worker(i, jobs, results, &wg)
}
go func() {
for i := int64(1); i <= 6; i++ {
jobs <- Task{ID: i, Title: fmt.Sprintf("到期任务%d", i)}
}
close(jobs) // 投递完毕,通知所有 worker 退出
}()
go func() {
wg.Wait() // 等所有 worker 退出
close(results) // 再关结果 channel
}()
for r := range results {
fmt.Println(r)
}
实测输出(顺序每次不同):
worker-1 关闭任务 #1(到期任务1)
worker-2 关闭任务 #2(到期任务2)
worker-3 关闭任务 #3(到期任务3)
worker-2 关闭任务 #5(到期任务5)
worker-1 关闭任务 #4(到期任务4)
worker-3 关闭任务 #6(到期任务6)
注意这里出现了两个关闭点,它们的顺序不能颠倒:
close(jobs)由生产者做,表示「没有新任务了」。close(results)必须等wg.Wait()返回之后才能做。如果提前关,worker 还在往里写就会 panic(send on closed channel)。
这就是 10.2 节说的「多个发送方时,由一个协调 goroutine 在所有人都结束后关闭」——这里的协调者就是那个 go func(){ wg.Wait(); close(results) }()。
10.3.3 worker 数量怎么定
没有万能公式,但有几条经验:
| 任务类型 | 建议 worker 数 | 理由 |
|---|---|---|
| CPU 密集(计算、压缩) | runtime.GOMAXPROCS(0) | 超过核数只会增加切换开销 |
| IO 密集(数据库、HTTP) | 核数的 2~10 倍,按下游限流反推 | 大部分时间在等,可以多开 |
| 有硬性配额的下游 | 直接取配额上限 | 比如「对方 API 限制 5 QPS」 |
| 不确定 | 先用固定值 + 压测调 | 别靠猜 |
关键是让 worker 数成为一个显式的配置项,而不是散落在代码里的魔法数字。第 15 章会讲怎么用 flag/env 把它做成可配置项。
另外,worker 数确定后,jobs channel 的缓冲容量也要一起考虑:缓冲为 0 时生产者和 worker 必须「碰面」交接,缓冲为 N 时可以预存 N 个任务。缓冲不是越大越好——它只是把压力从「生产者阻塞」换成了「内存占用」。
10.3.4 pipeline:把处理拆成阶段
worker pool 解决的是「一个阶段的并发」,pipeline 解决的是「多个阶段的串联」。
pipeline 的每个阶段都是一个函数:接收一个 channel,返回一个 channel。每个阶段内部起自己的 goroutine,用 defer close(out) 标记结束。
// 阶段一:生成
func gen(nums ...int) <-chan int {
out := make(chan int)
go func() {
defer close(out)
for _, n := range nums {
out <- n
}
}()
return out
}
// 阶段二:加工
func square(in <-chan int) <-chan int {
out := make(chan int)
go func() {
defer close(out)
for n := range in {
time.Sleep(10 * time.Millisecond) // 模拟耗时
out <- n * n
}
}()
return out
}
主流程读起来像一句自然语言:
for v := range square(gen(1, 2, 3, 4, 5)) {
fmt.Println("结果:", v)
}
实测输出:
结果: 1
结果: 4
结果: 9
结果: 16
结果: 25
串行 pipeline 耗时约 60ms
这个「每个阶段一个 goroutine」的写法有两个好处:
- 阶段可复用:
square不知道上游是谁,只知道自己从<-chan int读、往<-chan int写。它可以直接接到另一个gen上。 - 流式处理:不需要把所有数据先攒进一个切片。上游一边产、下游一边处理,内存占用与数据总量无关。
10.3.5 fan-out / fan-in
pipeline 的威力在于阶段内部可以横向扩展:
- fan-out:一个阶段启动多个 goroutine 消费同一个上游 channel。
- fan-in:把多个 goroutine 的输出汇进一个 channel。
10.3.2 的 worker pool 其实就是 fan-out + fan-in。下面把「每个阶段都并发」的版本拼出来:
func merge(cs ...<-chan string) <-chan string {
out := make(chan string)
var wg sync.WaitGroup
for _, c := range cs {
wg.Add(1)
go func(ch <-chan string) {
defer wg.Done()
for v := range ch {
out <- v
}
}(c)
}
go func() {
wg.Wait()
close(out) // 所有输入都枯竭后才关闭输出
}()
return out
}
merge 是 fan-in 的标准实现:它接收任意多个输入 channel,为每个输入起一个搬运 goroutine,最后用一个协调 goroutine 在全部输入关闭后关闭输出。
这里有一个必须注意的细节:for _, c := range cs 里的 c 在 Go 1.22 之后是每次迭代新变量,所以 go func(ch <-chan string) 其实可以不传参。但显式传参仍然更清晰,而且它在 for 循环体外声明变量的场景下是唯一正确的写法(10.1 节讲过)。
10.3.6 背压:让快的一方慢下来
pipeline 有一个很优雅的性质:channel 本身就是背压机制。
- 上游写得比下游读得快时,如果 channel 无缓冲或缓冲已满,上游的
out <- v会阻塞。 - 阻塞会沿着 pipeline 一路往上游传,直到最源头的生产者也被迫放慢。
这意味着你不需要写任何限速代码,只要不把 channel 缓冲开得过大,系统就会自动平衡。反过来,如果你为了「提高吞吐」把缓冲开到几万,就等于放弃了背压——上游会一路狂奔把内存吃光,直到 OOM。
所以缓冲容量的选择标准是:够吸收正常的抖动,不够吸收持续的过载。通常几十到几百的量级就足够。
10.3.7 把 worker pool 接进 TaskAPI
现在给 TaskAPI 实现「批量关闭到期任务」。这是本节所有概念的合体:固定 4 个 worker、无缓冲 jobs channel 做背压、atomic.Int64 统计处理数。
package main
import (
"fmt"
"sync"
"sync/atomic"
"time"
)
type Task struct {
ID int64
Title string
Done bool
}
func main() {
const total = 12
const workers = 4
jobs := make(chan Task)
var processed atomic.Int64
var wg sync.WaitGroup
for w := 1; w <= workers; w++ {
wg.Go(func() {
for t := range jobs {
time.Sleep(30 * time.Millisecond) // 模拟写库
processed.Add(1)
fmt.Printf("worker 关闭任务 #%d %s\n", t.ID, t.Title)
}
})
}
start := time.Now()
go func() {
defer close(jobs)
for i := int64(1); i <= total; i++ {
jobs <- Task{ID: i, Title: fmt.Sprintf("到期任务%d", i)}
}
}()
wg.Wait()
fmt.Printf("共处理 %d 个,耗时 %v(4 并发,单任务 30ms)\n",
processed.Load(), time.Since(start).Round(time.Millisecond))
}
实测输出:
worker 关闭任务 #2 到期任务2
worker 关闭任务 #1 到期任务1
worker 关闭任务 #4 到期任务4
worker 关闭任务 #3 到期任务3
...
共处理 12 个,耗时 93ms(4 并发,单任务 30ms)
12 个任务 ÷ 4 个 worker × 30 毫秒 = 90 毫秒,实测 93 毫秒,基本吻合。如果把 workers 改成 12,耗时就会降到 30 毫秒出头——这正是 10.3.1 表格里那组数字的来源。
注意这里用的是 wg.Go(Go 1.25 引入的 sync.WaitGroup 方法),它等价于 wg.Add(1) 加一个 defer wg.Done() 的 goroutine。第 11.2 节会完整讲它的语义和坑。
10.3.8 本节不写什么
并发是 Go 最容易写错、也最容易写深的领域。本节刻意停在「会用 + 能验证」这条线上,下面这些话题留给《Go 语言高级编程》卷:
errgroup:带错误传播和取消的 goroutine 组。semaphore:信号量式的并发限制。goleak:在测试里检测 goroutine 泄漏。- 结构化并发(structured concurrency)的整体范式。
- 调度器内部原理:GMP 三元组、work stealing、sysmon、抢占式调度。
- 垃圾回收的并发机制。
卷一的判断标准很简单:你写的每一个并发程序都应该能用 go run -race 跑通,并且你能说清「谁启动、谁等待、谁取消、谁关闭」。做到这一点,本章的目标就达到了。
10.3.9 小结与练习
- worker pool 解决「限并发」,pipeline 解决「分阶段」;两者可以叠加。
- worker 的退出条件是输入 channel 被关闭,
for range会自动结束。 - 多个发送方时,用一个
wg.Wait()+close(out)的协调 goroutine 统一关闭。 - channel 缓冲就是背压旋钮,开太大等于放弃背压。
- worker 数由下游承受力反推,不做魔法数字。
练习:
- 把 10.3.7 的
workers改成 1、4、12,记录耗时,画出「并发数—耗时」曲线,找出收益递减的拐点。 - 给 worker 加一个「处理失败就重试一次」的逻辑,用
atomic.Int64分别统计成功数和重试数。 - 把 10.3.4 的
gen、square扩成三阶段(生成 → 平方 → 求和),用merge让第二级并发 3 个 worker,观察输出顺序的变化。
下一章我们处理并发里最危险的部分:多个 goroutine 同时读写同一份数据时会发生什么,以及 -race 如何帮你抓住它。
阅读导航:上一节:10.2 channel 与 select · 下一节:11.1 Mutex/RWMutex 与 atomic 。
继续阅读
探索更多技术文章
浏览归档,发现更多关于系统设计、工具链和工程实践的内容。