Go 并发模式与 Goroutine 生命周期

# Go 并发模式与 Goroutine 生命周期

本篇目标

在基础 Goroutine、Channel 与 Context 之上,掌握生产代码最常见的 Worker Pool、Pipeline、背压、关闭顺序和泄漏排查。

明确 Goroutine 所有者实现有界并发传播取消避免泄漏与死锁

# 先记住一句话

每个启动的 Goroutine 都必须能回答 “谁让它停、谁等待它结束、阻塞时怎样取消” ;Channel 用于传递数据和同步所有权,Mutex 用于保护共享内存,两者按问题选择而不是互相替代。

# 启动 Goroutine 前先写退出条件

启动一个 Goroutine
├─ 输入从哪里来
├─ 输出由谁消费
├─ Channel 由谁关闭
├─ Context 何时取消
├─ 发生错误怎样上报
└─ 上层怎样等待结束
1
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
                队列已满
                    ↓
          等待空位,或者取消、超时
1
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 关联资源。
	})
}
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
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,调用方无法误发送或关闭。
}
1
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
}
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

只在接收输入时检查 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。
}
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

& 和 * 都与指针有关,但方向不同:&value 从值取得指针,*Type 表示指向该类型的指针;当 * 用在指针变量前时,表示读取或修改指针所指向的值。

count := 10

// &count 取得 count 的地址;pointer 是普通变量名,不是 Go 关键字。
pointer := &count // pointer 的类型是 *int。

// *pointer 根据地址读取它指向的值。
value := *pointer // value 等于 10。

// *pointer 放在赋值号左侧时,会修改它指向的原变量。
*pointer = 20 // count 现在等于 20。
1
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)
}
1
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,进一步掌握共享状态保护、并发等待和错误传播。

# 参考资料

上次更新时间: 2026年09月18日 02:14:27