预先创建固定数量的 worker 协程,让它们反复从任务通道中取任务执行,而不是每个任务都新建一个协程。这是 Go 并发编程中最常用的复用模式


裸写版

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
package main

import (
"fmt"
"sync"
"time"
)

func main() {
jobs := make(chan int, 3)
var wg sync.WaitGroup

// 启动 3 个 worker
for i := 0; i < 3; i++ {
wg.Add(1)
go func(id int) {
defer wg.Done()
for j := range jobs {
fmt.Printf("Worker %d job %d\n", id, j)
time.Sleep(time.Second)
}
}(i)
}

// 提交任务
for i := 0; i < 10; i++ {
jobs <- i
}
close(jobs) // 通知 worker 没有更多任务了
wg.Wait() // 等所有 worker 执行完毕再退出
}
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
package main

import (
"fmt"
"sync"
"time"
)

func main() {
jobs := make(chan int, 3)
var wg sync.WaitGroup

for i := 0; i < 3; i++ {
wg.Add(1)
go func(id int) {
defer wg.Done()
for j := range jobs {
fmt.Printf("Worker %d job %d\n", id, j)
time.Sleep(time.Second)
}
}(i)
}

for i := 0; i < 10; i++ {
jobs <- i
}
close(jobs)
wg.Wait()
}

三个核心步骤,没有封装:

  1. go func() { for j := range jobs { ... } } — 创建 worker,循环取任务
  2. jobs <- i — 提交任务到通道
  3. close(jobs) + wg.Wait() — 关闭通道 + 等待全部完成

精简版(封装结构体)

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
package main

import (
"fmt"
"sync"
"time"
)

type Pool struct {
jobs chan func() // 任务通道
wg sync.WaitGroup // 等待 worker 退出
}

func NewPool(sz int) *Pool {
return &Pool{jobs: make(chan func(), sz)}
}

func (p *Pool) Run(n int) {
for i := 0; i < n; i++ {
p.wg.Add(1)
go func(id int) {
defer p.wg.Done()
for job := range p.jobs {
fmt.Printf("w%d: 任务执行\n", id)
job()
}
}(i)
}
}

func (p *Pool) Sub(job func()) {
p.jobs <- job
}

func (p *Pool) Stop() {
close(p.jobs)
p.wg.Wait()
fmt.Println("已关闭")
}

func main() {
p := NewPool(3)
p.Run(3)
for i := 0; i < 10; i++ {
i := i
p.Sub(func() {
fmt.Printf("任务 %d 完成\n", i)
time.Sleep(time.Second)
})
}
p.Stop()
}
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
package main

import (
"fmt"
"sync"
"time"
)

type Pool struct {
jobs chan func()
wg sync.WaitGroup
}

func NewPool(sz int) *Pool {
return &Pool{jobs: make(chan func(), sz)}
}

func (p *Pool) Run(n int) {
for i := 0; i < n; i++ {
p.wg.Add(1)
go func(id int) {
defer p.wg.Done()
for job := range p.jobs {
fmt.Printf("w%d: 任务执行\n", id)
job()
}
}(i)
}
}

func (p *Pool) Sub(job func()) {
p.jobs <- job
}

func (p *Pool) Stop() {
close(p.jobs)
p.wg.Wait()
fmt.Println("已关闭")
}

func main() {
p := NewPool(3)
p.Run(3)
for i := 0; i < 10; i++ {
i := i
p.Sub(func() {
fmt.Printf("任务 %d 完成\n", i)
time.Sleep(time.Second)
})
}
p.Stop()
}

核心原理

预先创建固定数量的 worker 协程,让它们反复从 jobs 通道中取任务执行,而不是每个任务都新建一个协程。

核心就一行:for job := range p.jobs — worker 不退出,一直在循环取任务。

与信号量模式的区别

模式创建协程数关键代码
信号量10 个任务 → 创建 10 个协程(限制同时跑 3 个)sem <- 1; go func() { ...; <-sem }() 每次新建
Worker Pool10 个任务 → 创建 3 个 worker(反复复用)go func() { for job := range jobs { job() } }() 循环取任务

结构体与方法

WorkerPool 结构体

1
2
3
4
5
type WorkerPool struct {
jobs chan func() // 任务队列通道
wg sync.WaitGroup // 等待 worker 退出
size int // worker 数量
}
方法作用
NewWorkerPool(size)创建池,初始化 jobs 通道
Start()启动 N 个 worker,每个里面 for range jobs 循环取任务
Submit(job func())jobs 发送任务,通道满时自动阻塞限流
Stop()close(jobs) 关闭通道 + wg.Wait() 等待所有 worker 退出

执行流程

步骤 1 — 创建通道

jobs := make(chan func(), size),缓冲大小 = worker 数量。

步骤 2 — 启动 worker

启动 N 个 worker 协程,每个里面 for range jobs 循环取任务。

步骤 3 — 提交任务

主协程往 jobs 通道提交任务(Send),worker 自动取出执行(Receive)。

步骤 4 — 关闭通道

close(jobs) → worker 的 for range 循环结束,协程自然退出。

步骤 5 — 等待退出

wg.Wait() 确认所有 worker 都已退出,优雅关闭。


详细步骤拆解

步骤一:创建通道

1
jobs: make(chan func(), size)
  • 创建一个带缓冲的通道,存放 func() 类型的任务函数
  • 缓冲大小 = worker 数量,确保每个 worker 都能有一个任务在排队

类比: 建了一个流水线工作台,上面放了 3 个托盘,每个 worker 对应一个托盘位置。

步骤二:启动 worker(关键:常驻协程)

1
2
3
4
5
6
7
8
9
10
11
func (p *WorkerPool) Start() {
for i := 0; i < p.size; i++ { // p.size = 3,循环 3 次
p.wg.Add(1)
go func(id int) {
defer p.wg.Done()
for job := range p.jobs { // ← 循环不结束,协程就一直活着
job()
}
}(i)
}
}

从哪里体现 3 个常驻协程?

  1. for i := 0; i < p.size; i++ — 循环 3 次,每次 go func(...) 创建一个协程,共 3 个
  2. for job := range p.jobs — 这个循环只有在通道关闭时才结束,所以协程不会退出,一直活着等待新任务

对比信号量模式:

1
2
3
4
// 信号量模式:每次循环都 go 一个新的,10 次 = 10 个协程
for i := 0; i < 10; i++ {
go func() { ... }()
}
1
2
// 协程池:只在 Start 时创建 3 个,后面不再新建
Start() 创建: worker1, worker2, worker3 ← 就这一次,后面复用

步骤三:提交任务

1
2
3
func (p *WorkerPool) Submit(job func()) {
p.jobs <- job // 通道满了会阻塞,自动限流
}

主协程往 jobs 通道发送任务函数。如果通道满了(所有 worker 都在忙且托盘都占满),Submit 会阻塞等待,自动实现限流。

步骤四:关闭通道

1
close(p.jobs)

关闭通道后,for job := range p.jobs 会收到通道关闭的信号,循环自然结束。这是 Go 中通知消费者”没有更多任务了”的标准做法。

⚠️ 只有发送方才能关闭通道,不能在 worker 中关闭,否则其他 worker 再往通道发任务会 panic。

步骤五:等待所有 worker 退出

1
p.wg.Wait()

每个 worker 在 for range 结束后会执行 defer p.wg.Done()Wait() 阻塞直到所有 worker 的 Done() 都调用完毕,确保优雅关闭。


执行时序图

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
主协程                            worker 协程(3个常驻)

├─ NewWorkerPool(3)
│ jobs = make(chan func(), 3)

├─ pool.Start()
│ ├─ go Worker 0 → for range jobs (等待任务...)
│ ├─ go Worker 1 → for range jobs (等待任务...)
│ └─ go Worker 2 → for range jobs (等待任务...)

├─ pool.Submit(任务0) ──────────────→ Worker 0 取到,开始执行
├─ pool.Submit(任务1) ──────────────→ Worker 1 取到,开始执行
├─ pool.Submit(任务2) ──────────────→ Worker 2 取到,开始执行
│ (通道缓冲区满,后续 Submit 阻塞等待)

│ Worker 0 完成,取任务3...
│ Worker 1 完成,取任务4...
│ Worker 2 完成,取任务5...
│ ...以此类推...

├─ pool.Stop()
│ ├─ close(jobs) → 所有 worker 的 for range 结束
│ └─ wg.Wait() → 等待 Worker 0,1,2 都执行 Done()

└─ fmt.Println("协程池已关闭")

图解

3 个 worker 复用执行 10 个任务

1
2
3
4
5
         ┌─ worker1 (常驻) ─ 取任务0 → 取任务3 → 取任务6 → 取任务9
jobs 通道├─ worker2 (常驻) ─ 取任务1 → 取任务4 → 取任务7
└─ worker3 (常驻) ─ 取任务2 → 取任务5 → 取任务8

就这 3 个协程,反复复用,不新建不销毁

Worker Pool vs 信号量模式

信号量模式

  • 每个任务创建 1 个协程
  • 10 个任务 = 10 个协程(只是限制同时跑 3 个)
  • 协程创建和销毁有开销
  • 适合任务数量少、执行快的场景

Worker Pool

  • 预先创建固定数量 worker
  • 10 个任务 = 3 个协程反复复用
  • 无额外创建开销,内存占用稳定
  • 适合任务数量大、需要控制资源消耗的场景

常见问题

信号量通道和 Worker Pool 通道有什么区别?

  • 信号量通道 chan struct{} — 缓冲类型,只传递空值,作用是限制并发数,不传递数据
  • Worker Pool 通道 chan func() — 缓冲类型,传递任务函数本身,作用是任务分发

信号量模式是”每个任务一个协程,用通道限流”,Worker Pool 是”固定数量协程,用通道分发任务”。

worker 数量设多大合适?

和信号量模式的容量设置原则一样:

  • CPU 密集型 — 设为 runtime.NumCPU()
  • I/O 密集型 — 可以设大一些(如 50~100),因为 I/O 等待时不占 CPU
  • 不确定 — 从 3~10 开始,根据实际压测调整

如果任务执行很慢,通道会爆吗?

不会爆。jobs 通道是带缓冲的,当缓冲区满了之后 Submit 会阻塞,主协程自动被限流。worker 处理完一个任务后释放一个缓冲区位置,主协程就能继续提交。

如果不想阻塞,可以用 select + default 实现非阻塞提交:

1
2
3
4
5
6
select {
case p.jobs <- job:
// 提交成功
default:
// 通道已满,放弃或记录日志
}

为什么用 for range 而不是 for {} + <-p.jobs?

for job := range p.jobs 是 Go 中消费通道的标准写法:

  • 每次从通道接收一个值赋给 job
  • 当通道关闭且缓冲区为空时,循环自动结束
  • for { job := <-p.jobs } 更简洁,且能正确处理通道关闭

如果用 for { job := <-p.jobs } 需要额外检查 ok 来判断通道是否关闭:

1
2
3
4
5
6
7
for {
job, ok := <-p.jobs
if !ok {
break // 通道已关闭
}
job()
}

for range 自动帮你做了这件事。


📖 系列下一篇: 双层 Channel 任务分发模式

Worker Pool 的任务是 func(),结果在闭包内部处理。如果 worker 需要把计算结果传回来给调用方,就需要第二条通道。下一篇的双层 Channel 在 Worker Pool 基础上加一条 res 结果通道,实现生产者 / worker / 消费者各司其职。