goroutine泄漏问题与为什么需要池化
Go语言创建goroutine的成本极低,但这不等于可以无限创建。一个goroutine初始栈2KB,业务代码中持有闭包引用后实际占用远超2KB。批量处理100万条数据时,直接go func为每条数据启动一个goroutine,内存瞬间飙升至数GB,调度器延迟急剧恶化,CPU在上下文切换中空转。更隐蔽的问题是goroutine泄漏——如果goroutine内部阻塞在channel或锁上,外层又没有退出机制,这些goroutine会永久驻留内存,随运行时间累积最终OOM。
池化方案限制并发goroutine数量,复用worker避免频繁创建销毁,配合context实现超时和取消,是生产环境的标配。
固定大小goroutine池的实现方案
经典实现用带缓冲channel作为任务队列,N个worker goroutine从队列消费任务:
type WorkerPool struct {
tasks chan Task
wg sync.WaitGroup
workers int
ctx context.Context
cancel context.CancelFunc
}
type Task struct {
ID int
Data interface{}
Fn func(context.Context, interface{}) error
}
func NewWorkerPool(workers, queueSize int) *WorkerPool {
ctx, cancel := context.WithCancel(context.Background())
p := &WorkerPool{
tasks: make(chan Task, queueSize),
workers: workers,
ctx: ctx,
cancel: cancel,
}
p.start()
return p
}
func (p *WorkerPool) start() {
for i := 0; i < p.workers; i++ {
p.wg.Add(1)
go func(workerID int) {
defer p.wg.Done()
for {
select {
case <-p.ctx.Done():
return
case task, ok := <-p.tasks:
if !ok {
return
}
task.Fn(p.ctx, task.Data)
}
}
}(i)
}
}
func (p *WorkerPool) Submit(task Task) error {
select {
case p.tasks <- task:
return nil
default:
return fmt.Errorf("task queue full")
}
}
func (p *WorkerPool) Shutdown() {
p.cancel()
close(p.tasks)
p.wg.Wait()
}
Submit的非阻塞模式(default分支)在队列满时立即返回错误,避免生产者阻塞。需要阻塞等待时去掉default分支,但必须配合context超时防止死锁。
context超时控制与级联取消机制
context是Go并发安全的取消信号传播机制。父context取消时,所有子context自动取消,形成级联取消树:
func ProcessBatch(ctx context.Context, items []Item) error {
// 整个批次最多执行30秒
batchCtx, batchCancel := context.WithTimeout(ctx, 30*time.Second)
defer batchCancel()
// 每个单独任务最多5秒
itemCtx, itemCancel := context.WithTimeout(batchCtx, 5*time.Second)
result := make(chan error, 1)
go func() {
result <- processItem(itemCtx, items[0])
itemCancel()
}()
select {
case err := <-result:
return err
case <-itemCtx.Done():
return fmt.Errorf("item timeout: %w", itemCtx.Err())
case <-batchCtx.Done():
return fmt.Errorf("batch timeout: %w", batchCtx.Err())
}
}
关键习惯:context.WithTimeout返回的cancel函数必须调用,否则资源泄漏。用defer cancel()放在创建后立即调用是最安全的模式,即使超时已触发cancel也不会报错。
动态伸缩的goroutine池设计
固定池在负载波动大的场景下不够灵活。动态池根据队列积压情况增减worker:队列积压超过阈值时创建新worker,空闲超时后自动退出。
type DynamicPool struct {
tasks chan Task
minWorkers int
maxWorkers int
activeCount atomic.Int32
idleTimeout time.Duration
ctx context.Context
cancel context.CancelFunc
mu sync.Mutex
}
func (p *DynamicPool) dispatch() {
for {
select {
case <-p.ctx.Done():
return
case task := <-p.tasks:
if p.activeCount.Load() < int32(p.maxWorkers) {
p.mu.Lock()
p.activeCount.Add(1)
p.mu.Unlock()
go p.runWorker(task)
continue
}
p.runTask(task)
}
}
}
func (p *DynamicPool) runWorker(firstTask Task) {
defer p.activeCount.Add(-1)
idleTimer := time.NewTimer(p.idleTimeout)
defer idleTimer.Stop()
p.runTask(firstTask)
for {
select {
case <-p.ctx.Done():
return
case <-idleTimer.C:
if p.activeCount.Load() > int32(p.minWorkers) {
return
}
idleTimer.Reset(p.idleTimeout)
case task := <-p.tasks:
p.runTask(task)
idleTimer.Reset(p.idleTimeout)
}
}
}
minWorkers保证基线处理能力,maxWorkers防止失控扩张。idleTimeout建议30秒到2分钟,过短导致worker频繁销毁重建,过长浪费资源。
errgroup并发错误收集与优雅退出
标准库golang.org/x/sync/errgroup封装了并发任务编排+错误传播。Group.Go启动的goroutine中任一返回error,Group.Wait立即返回该错误,context自动取消通知其他goroutine退出:
import "golang.org/x/sync/errgroup"
func FetchAll(ctx context.Context, urls []string) ([]*http.Response, error) {
g, gctx := errgroup.WithContext(ctx)
results := make([]*http.Response, len(urls))
// 限制并发数
g.SetLimit(10)
for i, url := range urls {
i, url := i, url
g.Go(func() error {
req, err := http.NewRequestWithContext(gctx, "GET", url, nil)
if err != nil {
return err
}
resp, err := http.DefaultClient.Do(req)
if err != nil {
return err
}
results[i] = resp
return nil
})
}
if err := g.Wait(); err != nil {
return nil, err
}
return results, nil
}
SetLimit内部使用信号量控制并发数,无需手动构建channel。gctx在任一goroutine返回error时自动cancel,其他正在执行的HTTP请求收到取消信号中断连接。这种模式在批量API调用、并行数据抓取场景下非常实用。
原创文章,作者:小编,如若转载,请注明出处:https://www.yunthe.com/go-yu-yan-bing-fa-mo-shi-goroutine-chi-she-ji-yu-context/