商城首页欢迎来到中国正版软件门户

您的位置: 首页 > 文章列表 > 编程开发 > Go 语言如何实现一个简单的协程池(Worker Pool)?

Go 语言如何实现一个简单的协程池(Worker Pool)?

  发布于2026-07-10 阅读(0)

扫一扫,手机访问

直接 go func()

Go 语言如何实现一个简单的协程池(Worker Pool)?

为什么不能直接用 go func() 启动大量协程?

核心思路很简单:预启动固定数量的长期运行 worker,通过一个共享 channel 接收任务,避免反复创建销毁 goroutine 的开销。就这么几件事需要注意:

  • worker 数量通常设为 CPU 核心数 × 2~4,或者根据 I/O 密集型程度微调
  • 任务 channel 必须带缓冲(比如 make(chan Task, 100)),否则提交任务时可能阻塞 producer,影响上游
  • 不要在 worker 内部直接 recover panic——吞掉错误后调用方完全不知情;正确的做法是在 task 函数内自行处理,或者让 panic 向上传导由上层统一捕获

说白了,协程池存在的根本原因就是:不限制并发数,迟早要出事。

如何用 channel + for-range 实现基础 worker pool?

最简实现只需要一个任务 channel 和一组常驻 goroutine。每个 worker 循环从 channel 读任务并执行,代码量极少:

type Task func()

type WorkerPool struct {
    tasks chan Task
}

func NewWorkerPool(workerCount int) *WorkerPool {
    return &WorkerPool{
        tasks: make(chan Task, 100), // 缓冲区防止 producer 阻塞
    }
}

func (wp *WorkerPool) Start() {
    for i := 0; i < workerCount; i++ {
        go func() {
            for task := range wp.tasks {
                task()
            }
        }()
    }
}

func (wp *WorkerPool) Submit(task Task) {
    wp.tasks <- task // 非阻塞提交(因为有缓冲)
}

这里有个细节必须说清楚:range wp.tasks 会在 channel 关闭后自动退出循环,所以 shutdown 时需要显式 close(wp.tasks)。但关闭后就不能再调用 Submit 了,否则会 panic:send on closed channel。这一点在实际编码中很容易踩坑。

如何安全地停止协程池并等待任务完成?

直接 close(wp.tasks) 只能阻止新任务进入,已经接收但尚未执行的任务仍然会跑完,而正在执行的 task 却无法中断。要真正“等待所有任务结束”,必须引入额外的同步机制。

  • sync.WaitGroup 记录活跃 task 数量:Submit 前 wg.Add(1),task 执行完后 wg.Done()
  • Stop 方法先 close channel,再 wg.Wait() 等待全部 task 返回
  • 千万不要在 worker 的 for 循环中调用 wg.Done()——因为 range 循环结束后 worker 就退出了,但 task 可能还在执行;wg.Done() 必须放在 task 函数内部

典型的错误就是把 wg.Done() 放在了 worker 循环末尾,导致 task 还没执行完 wg 就减了,Wait() 提前返回,造成部分任务丢失。这个坑几乎每个初学协程池的人都会遇到。

要不要加 context 控制单个 task 超时或取消?

当然要。基础版 pool 对单个 task 完全没有感知,一旦某个 task 卡住(比如网络 hang、死循环),整个 worker 就被占住,后续任务排队等待。解决方式是在 task 层面支持 context.Context

type Task func(ctx context.Context) error

// Submit 时传入带 timeout 的 ctx
func (wp *WorkerPool) Submit(task Task, timeout time.Duration) {
    ctx, cancel := context.WithTimeout(context.Background(), timeout)
    defer cancel()
    wp.tasks <- func() { task(ctx) } // 包一层适配
}

这样一来,每个 task 可以自己决定是否响应 cancel,worker 完全不需要感知上下文——既保持了 pool 的简洁,又赋予了 task 级别的可控性。

不过有个问题:如果 task 内部压根没检查 ctx.Done(),那么超时设置形同虚设。真正棘手的是那些不支持 context 的第三方函数(比如旧版本的 http.Get),这时只能靠启动子 goroutine + select + channel 超时中转,但会增加复杂度和内存开销。究竟值不值得这么做,得看业务场景。

本文转载于:https://www.php.cn/faq/2393102.html 如有侵犯,请联系zhengruancom@outlook.com删除。
免责声明:正软商城发布此文仅为传递信息,不代表正软商城认同其观点或证实其描述。

热门关注