大家好,我是公号「左诗右码」作者 Alex。
上一篇我们介绍了用 Semaphore(信号量)模式来简单控制并发。但在生产环境中,我们往往需要一个更健壮、可复用、可监控的 Goroutine Pool(协程池)。
市面上有很多优秀的协程池库(如 ants),但作为 Go 开发者,自己手撸一个简易版是理解并发模式的最好方式。
今天,我们就带大家从零开始,设计并实现一个支持任务提交、固定 Worker 数、优雅关闭的协程池。
1. 为什么要用协程池?
虽然 Goroutine 很轻量(2KB),但也不是免费的。
- 内存开销:100 万个 G 也要占用 2GB+ 内存。
- 调度开销:G 太多会导致运行时调度压力变大,GC 扫描时间变长。
- 资源复用:通过池化技术,复用固定的 G 来处理任务,可以减少创建和销毁的开销。
2. 核心架构设计
我们的协程池 Pool 需要包含:
- Task Queue:一个缓冲 Channel,存放待执行的任务。
- Workers:一组固定的 Goroutine,循环从 Channel 里取任务执行。
- Capacity:池子的大小(Worker 数量)。
- Shutdown:关闭机制。
2.1 定义结构体
type Task func() // 任务就是一个函数
type Pool struct {
capacity int // Worker 数量
taskChan chan Task // 任务队列
wg sync.WaitGroup // 用于等待所有 Worker 退出
quit chan struct{} // 关闭信号
}
2.2 初始化与启动
func NewPool(cap int, taskQueueLen int) *Pool {
return &Pool{
capacity: cap,
taskChan: make(chan Task, taskQueueLen),
quit: make(chan struct{}),
}
}
// 启动 Worker
func (p *Pool) Start() {
for i := 0; i < p.capacity; i++ {
p.wg.Add(1)
go p.worker(i)
}
}
// Worker 逻辑
func (p *Pool) worker(workerID int) {
defer p.wg.Done()
fmt.Printf("Worker %d started\n", workerID)
for {
select {
case task, ok := <-p.taskChan:
if !ok {
return // Channel 关闭,退出
}
task() // 执行任务
case <-p.quit:
return // 收到关闭信号,退出
}
}
}
2.3 提交任务
func (p *Pool) Submit(t Task) {
select {
case p.taskChan <- t:
// 任务成功入队
case <-p.quit:
// 池子已关闭,不再接收
fmt.Println("Pool is closed, task rejected")
}
}
2.4 优雅关闭
func (p *Pool) Stop() {
close(p.quit) // 通知所有 worker 退出(优先于 taskChan)
p.wg.Wait() // 等待所有 worker 真正停下来
fmt.Println("All workers stopped")
}
3. 进阶:如何处理 Panic?
如果用户提交的任务 (task()) 发生了 Panic,会导致对应的 Worker 协程崩溃。如果 Worker 挂完了,池子就废了。
我们需要在 Worker 中捕获 Panic。
优化后的 Worker:
func (p *Pool) worker(workerID int) {
defer p.wg.Done()
for {
select {
case task, ok := <-p.taskChan:
if !ok { return }
// 封装执行,捕获 Panic
func() {
defer func() {
if r := recover(); r != nil {
fmt.Printf("Worker %d panic: %v\n", workerID, r)
}
}()
task()
}()
case <-p.quit:
return
}
}
}
4. 实战演示
func main() {
// 创建一个容量为 3 的池子
pool := NewPool(3, 10)
pool.Start()
// 提交 10 个任务
for i := 0; i < 10; i++ {
taskID := i
pool.Submit(func() {
fmt.Printf("Running task %d\n", taskID)
time.Sleep(500 * time.Millisecond)
if taskID == 5 {
panic("oops") // 模拟 panic
}
})
}
// 等待一会儿看效果
time.Sleep(3 * time.Second)
pool.Stop()
}
5. 总结
- 池化思想:限制并发数,复用资源,防止系统过载。
- Channel 作用:作为任务队列,连接生产者(Submit)和消费者(Worker)。
-
健壮性:必须在 Worker 中
recover,防止单个任务搞垮整个池子。 - 优雅关闭:利用 Context 或 Channel 广播退出信号,并用 WaitGroup 等待清理完成。
掌握了协程池,你的 Go 并发编程能力就上了一个台阶。
但在高并发场景下,除了计算资源的池化,对象内存的复用也同样重要。
下一篇,我们来聊聊 Go 的性能优化神器——sync.Pool。