在Golang的并发编程中,Goroutine是轻量级线程,使用起来非常方便,但如果无限制地创建Goroutine,会导致调度器负担加重,甚至引发内存溢出问题。Goroutine池的核心思想就是提前创建固定数量或可动态调整数量的协程,让这些协程循环处理提交的任务,实现协程的复用,减少创建和销毁的开销。

Goroutine池的基础结构设计
一个基础的Goroutine池通常包含几个核心部分:任务队列、工作协程集合、池的配置参数、关闭信号通道。下面我们先定义池的基础结构:
package main
import (
"context"
"errors"
"sync"
"time"
)
// Task 定义任务类型,是一个无参数无返回值的函数
type Task func() error
// GoroutinePool 协程池结构体
type GoroutinePool struct {
capacity int // 池的最大容量,即最大工作协程数
workerNum int // 当前工作协程数
taskQueue chan Task // 任务队列,用于存放待处理的任务
ctx context.Context
cancel context.CancelFunc
wg sync.WaitGroup // 用于等待所有工作协程退出
lock sync.Mutex // 保护workerNum等字段的互斥锁
closed bool // 池是否已关闭的标记
}
初始化与启动Goroutine池
初始化池的时候需要设置容量和任务队列的长度,同时启动指定数量的工作协程,这些协程会循环从任务队列中取任务执行。
// NewGoroutinePool 创建新的协程池
// capacity: 池的最大工作协程数
// queueSize: 任务队列的长度
func NewGoroutinePool(capacity int, queueSize int) (*GoroutinePool, error) {
if capacity <= 0 {
return nil, errors.New("pool capacity must be greater than 0")
}
if queueSize < 0 {
queueSize = capacity // 默认队列长度和容量一致
}
ctx, cancel := context.WithCancel(context.Background())
pool := &GoroutinePool{
capacity: capacity,
taskQueue: make(chan Task, queueSize),
ctx: ctx,
cancel: cancel,
closed: false,
}
// 初始化时启动capacity个工作协程
pool.startWorkers()
return pool, nil
}
// startWorkers 启动工作协程
func (p *GoroutinePool) startWorkers() {
p.lock.Lock()
defer p.lock.Unlock()
for i := 0; i < p.capacity; i++ {
p.wg.Add(1)
p.workerNum++
go p.worker()
}
}
工作协程的实现逻辑
工作协程的核心逻辑是循环监听任务队列和上下文的取消信号,当收到任务时执行任务,收到取消信号时退出循环,完成协程的回收。
// worker 工作协程的主逻辑
func (p *GoroutinePool) worker() {
defer p.wg.Done()
for {
select {
case task, ok := <-p.taskQueue:
if !ok {
// 任务队列被关闭,退出协程
return
}
// 执行任务,捕获可能的panic避免协程退出
func() {
defer func() {
if r := recover(); r != nil {
// 这里可以记录panic日志,实际项目中可替换为日志组件
}
}()
_ = task()
}()
case <-p.ctx.Done():
// 收到上下文取消信号,退出协程
return
}
}
}
任务提交与池的关闭
任务提交需要往任务队列中写入任务,同时要处理池已关闭的情况。关闭池的时候需要先取消上下文,关闭任务队列,再等待所有工作协程退出。
// Submit 提交任务到协程池
func (p *GoroutinePool) Submit(task Task) error {
p.lock.Lock()
defer p.lock.Unlock()
if p.closed {
return errors.New("goroutine pool is closed")
}
select {
case p.taskQueue <- task:
return nil
case <-time.After(100 * time.Millisecond): // 提交超时,避免阻塞调用方
return errors.New("submit task timeout")
}
}
// Close 关闭协程池
func (p *GoroutinePool) Close() {
p.lock.Lock()
if p.closed {
p.lock.Unlock()
return
}
p.closed = true
p.lock.Unlock()
// 先取消上下文,通知工作协程退出
p.cancel()
// 关闭任务队列,避免新的任务写入
close(p.taskQueue)
// 等待所有工作协程退出
p.wg.Wait()
}
Goroutine池的资源优化方法
动态调整池容量
固定的池容量无法适配所有场景,我们可以根据任务队列的长度动态调整工作协程的数量,当队列任务积压时扩容,任务较少时缩容,避免资源浪费。
// AdjustCapacity 动态调整池的容量
func (p *GoroutinePool) AdjustCapacity(newCapacity int) error {
if newCapacity <= 0 {
return errors.New("new capacity must be greater than 0")
}
p.lock.Lock()
defer p.lock.Unlock()
if p.closed {
return errors.New("pool is closed")
}
oldCapacity := p.capacity
p.capacity = newCapacity
// 如果新容量大于旧容量,启动新的工作协程
if newCapacity > oldCapacity {
needAdd := newCapacity - oldCapacity
for i := 0; i < needAdd; i++ {
p.wg.Add(1)
p.workerNum++
go p.worker()
}
}
// 如果新容量小于旧容量,通过上下文通知多余协程退出(这里简化逻辑,实际可通过标记控制)
// 注意:缩容时需要保证已经在执行的任务不会中断,这里仅做容量标记修改,实际项目可优化协程退出逻辑
return nil
}
限制任务超时与异常控制
为了避免单个任务执行时间过长阻塞协程,可以给任务执行增加超时控制,同时统一捕获任务执行中的异常,防止单个任务的错误导致整个协程退出。
// SubmitWithTimeout 提交带超时的任务
func (p *GoroutinePool) SubmitWithTimeout(task Task, timeout time.Duration) error {
return p.Submit(func() error {
// 创建带超时的上下文
ctx, cancel := context.WithTimeout(context.Background(), timeout)
defer cancel()
done := make(chan error, 1)
go func() {
done <- task()
}()
select {
case err := <-done:
return err
case <-ctx.Done():
return errors.New("task execute timeout")
}
})
}
避免协程泄漏
协程泄漏是Goroutine池常见的问题,主要原因是工作协程没有正确退出,或者任务队列中一直有未处理的任务导致协程阻塞。我们的实现中通过上下文取消、关闭任务队列、等待所有协程退出的逻辑,已经基本避免了协程泄漏的问题,实际使用中还需要注意不要在任务中启动新的未受控的Goroutine。
使用示例
下面是一个简单的使用示例,展示如何创建池、提交任务、关闭池:
func main() {
// 创建容量为5,任务队列长度为10的协程池
pool, err := NewGoroutinePool(5, 10)
if err != nil {
panic(err)
}
// 提交10个任务
for i := 0; i < 10; i++ {
idx := i
err := pool.Submit(func() error {
// 模拟任务执行耗时
time.Sleep(100 * time.Millisecond)
// 实际项目中可替换为业务逻辑
return nil
})
if err != nil {
// 处理提交失败的情况
}
}
// 等待所有任务执行完成(实际项目中可根据业务逻辑调整等待逻辑)
time.Sleep(200 * time.Millisecond)
// 关闭协程池
pool.Close()
}
GolangGoroutine_pool资源优化协程复用修改时间:2026-07-20 11:15:35