Go 并发模式与 Goroutine 生命周期
# Go 并发模式与 Goroutine 生命周期
在基础 Goroutine、Channel 与 Context 之上,掌握生产代码最常见的 Worker Pool、Pipeline、背压、关闭顺序和泄漏排查。
# 先记住一句话
每个启动的 Goroutine 都必须能回答 “谁让它停、谁等待它结束、阻塞时怎样取消” ;Channel 用于传递数据和同步所有权,Mutex 用于保护共享内存,两者按问题选择而不是互相替代。
# 启动 Goroutine 前先写退出条件
启动一个 Goroutine
├─ 输入从哪里来
├─ 输出由谁消费
├─ Channel 由谁关闭
├─ Context 何时取消
├─ 发生错误怎样上报
└─ 上层怎样等待结束
2
3
4
5
6
7
如果这些问题没有答案,go func() 只是把同步问题变成更难定位的异步问题。
# 一个有背压的 Worker Pool
背压(Backpressure)是指:下游处理不过来时,系统反过来限制上游继续提交任务,避免任务无限堆积。例如,任务每秒产生 100 个,而 Worker 每秒只能处理 20 个;如果不限制提交速度,积压任务会持续占用内存,最终可能拖垮进程。
常见的背压策略有:
- 等待:队列有空位后再提交;
- 拒绝:队列已满就立即返回 “系统繁忙”;
- 超时:等待一段时间,仍无法提交就放弃;
- 限流:提前限制上游生产任务的速度;
- 丢弃或降级:对不重要的任务采用简化处理。
背压不是某个单独的 Go API,而是一种过载控制机制。
这个 Worker Pool 通过两项限制形成背压:
- Worker 数固定,限制同时处理的任务数量;
jobs是容量有限的 Channel,队列已满时Submit无法继续发送,只能等待空位,或者在 Context 取消、超时后返回错误。
任务提交方 → 有界 jobs Channel → 固定数量的 Worker
队列已满
↓
等待空位,或者取消、超时
2
3
4
下面的 Pool 有固定 Worker 数、有界任务队列、Context 取消和明确关闭顺序。它适合进程内短任务;需要持久化、跨实例重试的任务仍应使用外部队列。
// 文件位置:internal/worker/pool.go
package worker
import (
"context" // 向所有 Worker 广播 Pool 取消信号,也允许 Submit 响应调用方取消。
"errors" // 创建可由 errors.Is 判断的稳定错误值。
"sync" // 提供 WaitGroup、Once 和 RWMutex 等同步工具。
)
// 包级错误值可供调用方通过 errors.Is 稳定判断错误类别。
var ErrPoolClosed = errors.New("worker pool is closed")
type Job struct {
ID string // 任务唯一标识,便于日志和幂等处理。
Payload string // Worker 真正要处理的任务内容。
}
// Handler 定义 Worker 的处理函数签名;Context 用于取消,error 用于报告处理失败。
type Handler func(context.Context, Job) error
// Pool 持有任务队列、Worker 生命周期和关闭状态,是这些资源的所有者。
type Pool struct {
// Context 负责取消,WaitGroup 等待退出,Once 保证只关闭一次,Mutex 保护共享状态。
ctx context.Context // 所有 Worker 共同观察的 Pool 生命周期 Context。
cancel context.CancelFunc // 释放 ctx 关联资源并广播取消信号。
jobs chan Job // 有界任务队列,元素类型固定为 Job。
handle Handler // 每个 Worker 收到 Job 后调用的业务函数。
wg sync.WaitGroup // 等待已经启动的 Worker 全部退出。
once sync.Once // 保证 Close 的关闭流程最多执行一次。
mu sync.RWMutex // 让 Submit 与 Close 安全地读取或修改 closed 并操作 jobs。
closed bool // 标记是否已经开始关闭;零值 false 表示仍可提交。
}
// New 校验配置、创建 Pool,并立即启动固定数量的 Worker。
// size 是 Worker 数,queueSize 是等待队列容量,handle 是每个任务的处理函数。
func New(parent context.Context, size int, queueSize int, handle Handler) (*Pool, error) {
// 至少需要一个 Worker;queueSize 可以为 0,此时 jobs 是无缓冲 Channel。
if size < 1 || queueSize < 0 {
return nil, errors.New("size must be positive and queueSize cannot be negative")
}
if handle == nil {
return nil, errors.New("handler is required")
}
// Pool 的 Context 派生自 parent:上层取消会向下传播,Pool 也能主动调用 cancel。
ctx, cancel := context.WithCancel(parent)
// &Pool{...} 创建 Struct 并取得指针;未填写字段保留各自零值。
pool := &Pool{
ctx: ctx,
cancel: cancel,
jobs: make(chan Job, queueSize), // 有界队列把过载反馈给提交方。
handle: handle,
}
pool.wg.Add(size) // 启动前先登记 Worker 数,避免 Goroutine 过快退出导致计数时序错误。
for index := 0; index < size; index++ {
// index 只用于控制启动次数;每个 Worker 都执行同一个 run 方法。
go pool.run() // Pool 是这些 Goroutine 的所有者。
}
return pool, nil // 返回已经开始工作的 Pool;nil 表示构造成功。
}
// run 是单个 Worker 的主循环,持续等待 Pool 取消或新任务。
func (p *Pool) run() {
defer p.wg.Done() // Worker 退出时通知 Close。
for {
select {
case <-p.ctx.Done():
return // 进程级取消时停止正在等待新任务的 Worker。
case job, ok := <-p.jobs:
if !ok {
return // 队列关闭且已排空,Worker 正常退出。
}
// 调用业务处理函数;_ 表示教学示例暂时忽略返回的 error。
_ = p.handle(p.ctx, job)
// 生产代码必须把错误写入日志、指标或结果存储,不能静默丢弃。
}
}
}
// Submit 尝试把 Job 放入有界队列,并响应调用方或整个 Pool 的取消。
func (p *Pool) Submit(ctx context.Context, job Job) error {
p.mu.RLock() // 与 Close 串行,保证发送期间 Channel 不会被关闭。
defer p.mu.RUnlock()
if p.closed {
return ErrPoolClosed
}
if err := p.ctx.Err(); err != nil {
return err // Pool 已被上层取消时,不再接受可能无人消费的新任务。
}
// 队列已满时,发送分支暂时无法执行;Submit 会等待空位、调用方取消或 Pool 取消。
select {
case <-p.ctx.Done():
return p.ctx.Err()
case <-ctx.Done():
return ctx.Err() // 队列已满时,调用方仍能超时或取消。
case p.jobs <- job:
return nil // 成功进入队列;这不代表 Job 已经处理完成。
}
}
// Close 停止接收新任务,关闭队列,等待已入队任务处理完毕并释放 Context。
func (p *Pool) Close() {
// Once.Do 保证并发或重复调用 Close 时,关闭逻辑最多执行一次。
p.once.Do(func() {
p.mu.Lock()
p.closed = true // 先阻止新的提交,并等待进行中的 Submit 结束。
close(p.jobs) // 只有发送方 Pool 可以关闭任务 Channel。
p.mu.Unlock()
p.wg.Wait() // 等待 Worker 消费完已入队任务并退出。
p.cancel() // 释放 Context 关联资源。
})
}
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
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
这个教学实现有一个工程边界:任务开始执行后无法被单独取消,因为队列里只保存 Job。真实项目可以让 Job 带自己的 Context,但不能把会跨请求长期存活的任务绑定到短请求 Context;更常见的是保存 deadline、task ID 和持久化取消状态。
# 谁发送,谁关闭
关闭 Channel 表示 “以后不会再有值” ,不是 “清空” 或 “销毁” 。通常由唯一发送方或协调所有发送方的所有者关闭;接收方不应擅自关闭,因为它不知道是否仍有发送者。
// 文件位置:examples/channel_owner.go
package channelowner
func produce(values []int) <-chan int {
// <-chan int 是只接收 Channel 类型,限制调用方只能读取。
output := make(chan int) // 无缓冲 Channel 会让生产速度自然受消费速度限制。
// 单独的生产 Goroutine 逐个发送数据,使调用方可以边接收边处理。
go func() {
defer close(output) // 创建并发送 output 的 Goroutine 负责关闭。
for _, value := range values {
output <- value // 没有接收方时在这里阻塞,形成同步交接。
}
}()
return output // 返回只读 Channel,调用方无法误发送或关闭。
}
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
向已关闭 Channel 发送会 panic;从已关闭且已取空的 Channel 接收会立即得到零值。需要区分真实零值和关闭状态时使用 value, ok := <-channel。
# Pipeline 必须支持下游提前退出
// 文件位置:internal/pipeline/map.go
package pipeline
import "context" // 让 Pipeline 阶段能够响应上游或下游发出的取消信号。
// Map 从 input 接收整数,用 transform 转换后发送到返回的只读 Channel。
func Map(ctx context.Context, input <-chan int, transform func(int) int) <-chan int {
// transform 是作为参数传入的函数,用于替换每个元素的转换规则。
output := make(chan int) // 当前阶段拥有 output,因此也负责最终关闭它。
// Pipeline 阶段在独立 Goroutine 中运行,调用 Map 后可以立即开始消费 output。
go func() {
defer close(output)
for {
select {
case <-ctx.Done():
return // 下游取消后停止等待输入。
case value, ok := <-input:
if !ok {
return // 上游结束后关闭当前阶段输出。
}
mapped := transform(value) // 转换发生在当前 Pipeline Goroutine 中。
select {
case output <- mapped:
// 下游成功接收后继续下一项。
case <-ctx.Done():
return // 下游不再接收时,避免永久阻塞发送。
}
}
}
}()
return output
}
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
只在接收输入时检查 Context 不够:如果下游停止接收,Goroutine 可能永远卡在输出发送。因此每一个潜在阻塞点都要考虑取消路径。
# Mutex 和 Channel 怎样选择
| 问题 | 更自然的工具 |
|---|---|
| 保护一个计数器、Map 或缓存状态 | sync.Mutex / sync.RWMutex |
| 把任务所有权从生产者交给消费者 | Channel |
| 等待一组 Goroutine 完成 | sync.WaitGroup |
| 广播请求取消和截止时间 | context.Context |
| 只初始化一次共享资源 | sync.Once |
| 限制同时进入某段逻辑的数量 | 有缓冲 Channel 或专门信号量 |
不要为了套 “通过通信共享内存” 把简单计数器交给专门 Goroutine,也不要为了少写 Channel 而让多个 Goroutine 随意共享复杂状态。选择能最直接表达所有权和不变量的工具。
// 文件位置:internal/cache/cache.go
package cache
import "sync" // RWMutex 允许多个读取者并发,但写入仍然独占。
type Cache struct {
mu sync.RWMutex // 保护 items 的所有读写;Mutex 与 Map 放在同一个 Struct 中。
items map[string]string // 真正保存缓存 Key 和 Value 的共享 Map。
}
// *Cache 表示 “指向 Cache 的指针”,这里声明 New 返回的类型。
func New() *Cache {
// Cache{...} 创建 Cache 值;& 取得它的地址,因此结果类型是 *Cache。
return &Cache{items: make(map[string]string)}
}
// (c *Cache) 表示 c 是指针接收者,方法操作的是原来的 Cache,而不是它的副本。
func (c *Cache) Get(key string) (string, bool) {
c.mu.RLock() // 多个纯读取可以并发执行。
defer c.mu.RUnlock()
value, ok := c.items[key]
return value, ok // ok 用于区分 Key 不存在和 Value 恰好是空字符串。
}
func (c *Cache) Set(key string, value string) {
c.mu.Lock() // Map 写入和其他读写互斥。
defer c.mu.Unlock()
c.items[key] = value // 只有持有写锁时才能修改普通 Map。
}
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
& 和 * 都与指针有关,但方向不同:&value 从值取得指针,*Type 表示指向该类型的指针;当 * 用在指针变量前时,表示读取或修改指针所指向的值。
count := 10
// &count 取得 count 的地址;pointer 是普通变量名,不是 Go 关键字。
pointer := &count // pointer 的类型是 *int。
// *pointer 根据地址读取它指向的值。
value := *pointer // value 等于 10。
// *pointer 放在赋值号左侧时,会修改它指向的原变量。
*pointer = 20 // count 现在等于 20。
2
3
4
5
6
7
8
9
10
可以简单记成:& 是从值取得指针,* 是声明指针类型,或通过指针访问原值。
# Context 只携带请求级信息
Context 适合截止时间、取消信号和跨 API 边界的请求级元数据。它不应:
- 存在长期 Struct 字段里;
- 作为可选参数传
nil; - 用来传数据库连接或普通业务参数;
- 被库函数无故替换成
context.Background(),从而丢失上游取消。
// 文件位置:internal/note/service.go
package note
import "context" // Repository 调用使用 Context 接收超时和取消信号。
// Repository 只声明 Service 读取 Note 所需的最小能力。
type Repository interface {
Get(ctx context.Context, id string) (Note, error)
}
type Note struct {
ID string // 笔记唯一标识。
Title string // 笔记标题。
}
type Service struct {
repository Repository // 具体实现由上层创建后注入。
}
// Get 不创建新的 Background Context,而是把调用方的 ctx 原样传给 Repository。
func (s *Service) Get(ctx context.Context, id string) (Note, error) {
// 原样向下传递调用方 Context,使数据库能观察断连和超时。
return s.repository.Get(ctx, id)
}
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
# Goroutine 泄漏怎样出现
典型泄漏不是 “忘记释放 Goroutine 对象” ,而是 Goroutine 永久阻塞且仍被运行时保留:
- 向无人接收的 Channel 发送;
- 从永远不会关闭或发送的 Channel 接收;
- 下游提前返回,上游不知道取消;
- Ticker 不停止,循环没有退出条件;
- 网络调用没有 Context 或超时;
- 锁顺序不一致造成死锁。
排查时结合 Goroutine Profile、阻塞 Profile、Trace 和测试。最新 Go 工具可能提供更多泄漏诊断能力,但代码设计仍要先保证生命周期明确。
# 高频面试题与回答
1. Channel 应该由谁关闭?参考答案
通常由确定不会再发送数据的一方关闭,也就是唯一发送者或协调所有发送者的所有者。接收方通常不能判断其他发送者是否结束,擅自关闭可能让仍在发送的 Goroutine panic。
2. 怎样避免 Goroutine 泄漏?参考答案
启动时就定义退出条件、取消来源和等待者;所有可能阻塞的 Channel 与 I/O 操作都支持 Context;限制队列和并发;发送方负责关闭;进程关闭时先停止接收新工作,再取消或排空任务并等待 Goroutine 结束。
3. Mutex 与 Channel 应怎样选择?参考答案
保护一个共享状态或短临界区时 Mutex 更直接;传递数据所有权、建立阶段流水线或协调生产消费时 Channel 更自然。两者都不是越多越好,关键是让所有权和不变量容易推理。
4. select 中多个分支同时就绪时怎样选择?参考答案
Go 会从已经就绪的通信分支中伪随机选择一个,不保证源码顺序和业务优先级;没有分支就绪时会阻塞,除非存在 default。因此不能用分支排列实现严格优先级,滥用 default 还可能形成忙循环并持续占用 CPU。
5. 数据竞争、死锁和 Goroutine 泄漏有什么区别?参考答案
数据竞争是多个 Goroutine 未同步访问同一内存且至少一方写入;死锁是相关执行单元都在等待,系统无法继续;Goroutine 泄漏是任务失去用途却一直没有退出。它们可能同时出现,但排查手段不同,分别重点使用 Race Detector、阻塞或 Goroutine Profile,以及生命周期和取消链路检查。
# 接下来学什么
下一篇学习 Go 同步原语、atomic 与 errgroup,进一步掌握共享状态保护、并发等待和错误传播。