《Go 语言编程入门》10.3 worker pool 与 pipeline

一个消费者太慢,无限制起 goroutine 又会打爆下游。本节用固定数量的 worker pool 消费任务队列,把「生成 → 加工 → 汇总」串成 pipeline,实测 4 并发与无限制并发的耗时差异,并给出 TaskAPI 批量关闭到期任务的完整实现。

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):

写法实测耗时峰值并发
每个任务一个 goroutine31 ms12
4 个 worker 消费队列93 ms4
无限制并发: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)
	}
}

三个要点:

  1. worker 的退出条件是 jobs 被关闭。for range jobs 会在 channel 关闭且缓冲读空后自动结束,不需要额外的退出信号。
  2. worker 自己不去关 jobs。它只是接收方,关闭是生产者的事(10.2 节的铁律)。
  3. 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 小结与练习

  1. worker pool 解决「限并发」,pipeline 解决「分阶段」;两者可以叠加。
  2. worker 的退出条件是输入 channel 被关闭,for range 会自动结束。
  3. 多个发送方时,用一个 wg.Wait() + close(out) 的协调 goroutine 统一关闭。
  4. channel 缓冲就是背压旋钮,开太大等于放弃背压。
  5. 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 。

继续阅读

探索更多技术文章

浏览归档,发现更多关于系统设计、工具链和工程实践的内容。

全部文章 返回首页

「golang」更多文章

  1. 《Go 语言编程实战》目录
  2. 《Go 语言编程实战》18.3 上线、观测与迭代
  3. 《Go 语言编程实战》18.2 故障演练