在Go语言开发中,频繁创建goroutine可能导致系统调度开销过大甚至内存溢出。构建一个可扩展的协程池框架,可以有效控制并发规模,复用执行资源,并方便后续功能扩展。

协程池的基本组成
一个基础的协程池通常包含任务队列、worker集合以及调度逻辑。任务以函数形式提交到队列,空闲worker从队列中取任务执行。
核心结构体设计
我们可以用以下结构描述协程池:
package main
import (
"sync"
)
// Task 代表一个可执行任务
type Task func()
// Pool 协程池结构
type Pool struct {
tasks chan Task
wg sync.WaitGroup
workers int
}
// NewPool 创建指定worker数量的协程池
func NewPool(workerNum int) *Pool {
return &Pool{
tasks: make(chan Task, 100),
workers: workerNum,
}
}
// Run 启动worker
func (p *Pool) Run() {
for i := 0; i < p.workers; i++ {
p.wg.Add(1)
go func() {
defer p.wg.Done()
for task := range p.tasks {
task()
}
}()
}
}
// Submit 提交任务
func (p *Pool) Submit(t Task) {
p.tasks <- t
}
// Stop 关闭池并等待完成
func (p *Pool) Stop() {
close(p.tasks)
p.wg.Wait()
}
实现可扩展性
上述代码是固定worker数量的池。要支持可扩展,可以引入动态扩容机制:当任务队列持续满时,临时增加worker;空闲一段时间后再回收。
动态扩容思路
- 监控任务队列长度与worker繁忙率
- 超过阈值时调用go启动新worker并计入管理器
- 使用timer检测空闲worker并安全退出
任务超时与错误处理
在并发框架中,任务可能panic或执行过久。可通过recover()捕获异常,并结合context控制超时。
func safeTask(t Task) {
defer func() {
if r := recover(); r != nil {
// 处理panic,避免worker退出
}
}()
t()
}
使用示例
下面演示如何提交多个任务并等待完成:
func main() {
pool := NewPool(5)
pool.Run()
for i := 0; i < 10; i++ {
idx := i
pool.Submit(func() {
println("task", idx)
})
}
pool.Stop()
}
通过上述方式,我们实现了一个简单但具备扩展基础的Golang协程池框架。实际业务中可继续加入优先级队列、指标统计等模块。