第 23 篇:实战并发(二):如何优雅地实现一个 Goroutine 池?

大家好,我是公号「左诗右码」作者 Alex。

上一篇我们介绍了用 Semaphore(信号量)模式来简单控制并发。但在生产环境中,我们往往需要一个更健壮、可复用、可监控的 Goroutine Pool(协程池)。

市面上有很多优秀的协程池库(如 ants),但作为 Go 开发者,自己手撸一个简易版是理解并发模式的最好方式。

今天,我们就带大家从零开始,设计并实现一个支持任务提交、固定 Worker 数、优雅关闭的协程池。


1. 为什么要用协程池?

虽然 Goroutine 很轻量(2KB),但也不是免费的。

  1. 内存开销:100 万个 G 也要占用 2GB+ 内存。
  2. 调度开销:G 太多会导致运行时调度压力变大,GC 扫描时间变长。
  3. 资源复用:通过池化技术,复用固定的 G 来处理任务,可以减少创建和销毁的开销。

2. 核心架构设计

我们的协程池 Pool 需要包含:

  1. Task Queue:一个缓冲 Channel,存放待执行的任务。
  2. Workers:一组固定的 Goroutine,循环从 Channel 里取任务执行。
  3. Capacity:池子的大小(Worker 数量)。
  4. 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. 总结

  1. 池化思想:限制并发数,复用资源,防止系统过载。
  2. Channel 作用:作为任务队列,连接生产者(Submit)和消费者(Worker)。
  3. 健壮性:必须在 Worker 中 recover,防止单个任务搞垮整个池子。
  4. 优雅关闭:利用 Context 或 Channel 广播退出信号,并用 WaitGroup 等待清理完成。

掌握了协程池,你的 Go 并发编程能力就上了一个台阶。

但在高并发场景下,除了计算资源的池化,对象内存的复用也同样重要。
下一篇,我们来聊聊 Go 的性能优化神器——sync.Pool

©著作权归作者所有,转载或内容合作请联系作者
【社区内容提示】社区部分内容疑似由AI辅助生成,浏览时请结合常识与多方信息审慎甄别。
平台声明:文章内容(如有图片或视频亦包括在内)由作者上传并发布,文章内容仅代表作者本人观点,简书系信息发布平台,仅提供信息存储服务。

友情链接更多精彩内容