vault backup: 2026-06-07 12:14:39

This commit is contained in:
2026-06-07 12:14:39 +08:00
parent 1a72bc4d82
commit 7ce9b83218
61 changed files with 8409 additions and 14429 deletions
+152 -149
View File
@@ -1,194 +1,197 @@
---
tags:
- Go
- golang
- go进阶
- 协程池
- 并发
tags: [go, golang, 协程池, 并发, WorkerPool]
create time: 2026-06-07 15:00
---
# 协程池
Go语言虽然有着高效的GMP调度模型,理论上支持成千上万的`goroutine`,但是`goroutine`过多,对调度,gc以及系统内存都会造成压力,这样会使我们的服务性能不升反降。常用做法可以用池化技术,构造一个协程池,把进程中的协程控制在一定的数量,防止系统中`goroutine`过多,影响服务性能。
## 协程池模型
协程池简单理解就是有一个池子一样的东西,里面装这个固定数量的`goroutine`,当有一个任务到来的时候,会将这个任务交给池子里的一个空闲的`goroutine`去处理,如果池子里没有空闲的`goroutine`了,任务就会阻塞等待。所以协程池有三个角色`Worker`,`Task`,`Pool`。
## 概述
### 属性定义
- `Worker`:用于执行任务的`goroutine`
- `Task`: 具体的任务
- `Pool`: 池子
虽然 Go 可以轻松创建数十万个 goroutine,但无限制地创建反而会导致调度开销和 GC 压力。协程池通过将活跃 goroutine 数量控制在合理范围内,实现性能与资源的平衡。
下面看一下各个角色的定义:
## 正文
#### Task定义
`Task`有一个函数成员,表示这个task具体的执行逻辑:
### 为什么需要协程池?
```go
type Task struct {
f func() error // 具体的执行逻辑
}
> [!question] 💭 思考
> Go 说可以开 10 万个 goroutine,那是不是意味着我应该每次都 `go func()`?
理论上可以,但实际上:
- **过多 goroutine** → GMP 调度器负担加重、CPU 上下文切换频繁
- **过多 goroutine** → GC 扫描对象增多、STW 时间变长
- **过多 goroutine** → 系统内存压力增大
协程池的核心思想:**控制并发度,而非消除并发**。
```mermaid
flowchart TD
A["任务队列<br/>Job Channel"] -->|"取任务"| B["Worker 1"]
A -->|"取任务"| C["Worker 2"]
A -->|"取任务"| D["Worker N"]
E["AddTask<br/>提交任务"] --> A
subgraph Pool["协程池 (固定 N 个 worker)"]
B
C
D
end
style A fill:#fff3e0
style E fill:#e8f5e9
style Pool fill:#e3f2fd
```
#### Pool定义
`Pool`有两个成员,`Capacity`表示池子里的worker的数量,即工作的`goroutine`的数量,`JobCh`表示任务队列用于存放任务,`goroutine`从这个`JobCh`获取任务执行任务逻辑:
```go
type Pool struct {
RunningWorkers int64
Capacity int64 // goroutine数量
JobCh chan *Task // 用于worker取任务
sync.Mutex
}
```
### 核心组件
#### Worker定义
```go
// p为Pool对象指针
for task := range p.JobCh {
do ...
}
```
执行任务单元,简单理解就是干活的`goroutine`,这个worker其实只做一件事情,就是不断的从任务队列里面取任务执行,而worker的数量就是协程池里协程的数量,由`Pool`的参数`WorkerNum`指定。
| 角色 | 说明 |
|------|------|
| Task | 封装待执行的业务逻辑(函数 + 参数) |
| Worker | 固定的 goroutine,从任务队列循环取任务执行 |
| Pool | 管理 Worker 数量和任务队列的容器 |
### 方法定义
```go
func NewTask(funcArg func() error) *Task
```
`NewTask`用于创建一个任务,参数是一个函数,返回值是一个`Task`类型。
### 最小实现
```go
func NewPool(Capacity int, taskNum int) *Pool
```
`NewPool`返回一个协程数量固定为`workerNum`协程池对象指针,其任务队列的长度为`taskNum`。
接下来主要介绍协程池的各个方法:
```go
func (p *Pool) AddTask(task *Task)
```
`AddTask`方法是往协程池添加任务,如果当前运行着的worker数量小于协程池worker容量,则立即启动一个协程worker来处理任务,否则将任务添加到任务队列。
```go
func (p *Pool) Run()
```
将协程池跑起来,启动一个worker来处理任务。
协程池处理任务流程图:
![协程池流程](https://golangstar.cn/assets/img/go语言系列/协程池/协程池1.png)
### 协程池实现
```go
package main
import (
"fmt"
"sync"
"sync/atomic"
"time"
"fmt"
"sync"
"sync/atomic"
)
// Task 封装一个可执行的任务
type Task struct {
f func() error // 具体的任务逻辑
f func() error
}
func NewTask(funcArg func() error) *Task {
return &Task{
f: funcArg,
}
func NewTask(f func() error) *Task {
return &Task{f: f}
}
// Pool 协程池
type Pool struct {
RunningWorkers int64 // 运行着的worker数量
Capacity int64 // 协程池worker容量
JobCh chan *Task // 用于worker取任务
sync.Mutex
capacity int // worker 数量
taskQueue chan *Task // 任务缓冲队列
wg sync.WaitGroup // 等待所有 worker 结束
running atomic.Int64 // 当前运行数
}
func NewPool(capacity int64, taskNum int) *Pool {
return &Pool{
Capacity: capacity,
JobCh: make(chan *Task, taskNum),
}
func NewPool(capacity int, queueSize int) *Pool {
return &Pool{
capacity: capacity,
taskQueue: make(chan *Task, queueSize),
}
}
func (p *Pool) GetCap() int64 {
return p.Capacity
// Start 启动 pool
func (p *Pool) Start() {
for i := 0; i < p.capacity; i++ {
p.wg.Add(1)
go p.worker(i)
}
}
func (p *Pool) incRunning() { // runningWorkers + 1
atomic.AddInt64(&p.RunningWorkers, 1)
func (p *Pool) worker(id int) {
defer p.wg.Done()
for task := range p.taskQueue {
if task != nil && task.f != nil {
task.f()
}
}
}
func (p *Pool) decRunning() { // runningWorkers - 1
atomic.AddInt64(&p.RunningWorkers, -1)
// Submit 提交任务
func (p *Pool) Submit(task *Task) bool {
select {
case p.taskQueue <- task:
return true
default:
return false // 队列已满,非阻塞拒绝
}
}
func (p *Pool) GetRunningWorkers() int64 {
return atomic.LoadInt64(&p.RunningWorkers)
// Stop 停止 pool,关闭任务队列等待所有 worker 完成
func (p *Pool) Stop() {
close(p.taskQueue)
p.wg.Wait()
}
```
func (p *Pool) run() {
p.incRunning()
go func() {
defer func() {
p.decRunning()
}()
for task := range p.JobCh {
task.f()
}
}()
}
// AddTask 往协程池添加任务
func (p *Pool) AddTask(task *Task) {
// 加锁防止启动多个 worker
p.Lock()
defer p.Unlock()
if p.GetRunningWorkers() < p.GetCap() { // 如果任务池满, 则不再创建 worker
// 创建启动一个 worker
p.run()
}
// 将任务推入队列, 等待消费
p.JobCh <- task
}
使用示例:
```go
func main() {
// 创建任务池
pool := NewPool(3, 10)
pool := NewPool(3, 10) // 3 个 worker,10 个任务缓冲
pool.Start()
for i := 0; i < 20; i++ {
// 任务放入池中
pool.AddTask(NewTask(func() error {
fmt.Printf("I am Task\n")
return nil
}))
}
for i := 0; i < 20; i++ {
pool.Submit(NewTask(func() error {
fmt.Printf("处理任务 %d\n", i)
return nil
}))
}
time.Sleep(1e9) // 等待执行
pool.Stop() // 优雅关闭
}
```
运行结果:
### 队列满了怎么办?
> [!question] 💭 思考
> 当任务生产速度超过消费速度,且队列也满了,你的协程池应该如何应对?
三种策略各有取舍:
| 策略 | 实现方式 | 优点 | 缺点 |
|------|---------|------|------|
| **阻塞等待** | `p.taskQueue <- task` | 不丢任务 | 生产者可能被卡住 |
| **非阻塞拒绝** | `select + default` | 响应快 | 可能丢任务 |
| **动态扩缩容** | 监控队列长度增减 worker | 自适应负载 | 实现复杂 |
> [!tip] 💡 推荐方案
> 大多数场景下,"有界队列 + 拒绝策略"是最实用的组合。拒绝时可以选择丢弃、记录日志告警,或降级处理。
### 优雅关闭
> [!warning] ⚠️ 协程池关闭的关键点
> 关闭时必须确保:1) 不再接收新任务;2) 已有任务全部执行完毕;3) 所有 worker goroutine 都退出。
```go
// 正确关闭流程
func (p *Pool) Shutdown() {
close(p.taskQueue) // 1. 关闭队列——worker 收到 range 退出信号
p.wg.Wait() // 2. 等待所有 worker 完成剩余任务
}
```
I am Task
I am Task
I am Task
I am Task
I am Task
I am Task
I am Task
I am Task
I am Task
I am Task
I am Task
I am Task
I am Task
I am Task
I am Task
I am Task
I am Task
I am Task
I am Task
I am Task
```
程序创建了一个`WorkerNum`为3,任务队列长度为10的协程池,往里面添加了20个任务,可以看到输出,一直只有3个`worker`在做任务,起到了控制`goroutine`数量的作用。
> [!note] 📝 不要直接调用 runtime.Goexit() 或 panic 来终止 goroutine
> 这会导致 defer 链不完整、资源泄漏等问题。始终用 channel close + WaitGroup 的方式优雅退出。
### 实战建议
> [!tip] 💡 协程池调参指南
> - **CPU 密集型**:worker 数量 ≈ CPU 核心数
> - **IO 密集型**:worker 数量 = CPU 核心数 × (1 + 等待时间/计算时间),通常 10~100 倍
> - **混合场景**:根据压测结果调整,观察 CPU 利用率和延迟 P99
> [!warning] ⚠️ 避免在 worker 中 panic
> 单个 worker panic 会导致整个程序崩溃。应在 worker 中使用 recover:
> ```go
> func (p *Pool) worker(id int) {
> defer func() {
> if r := recover(); r != nil {
> log.Printf("worker %d recovered from panic: %v", id, r)
> }
> p.wg.Done()
> }()
> for task := range p.taskQueue { ... }
> }
> ```
## 关联笔记
- [[hzh/GolangStar/Go语言进阶/Goroutine]]
- [[hzh/GolangStar/Go语言进阶/Sync]]
- [[hzh/GolangStar/Go语言进阶/Channel]]