扇出/扇入工作池上下文取消后Goroutine泄漏问题排查与修复
问题描述
实现了一个扇出/扇入工作池,预期上下文取消时所有worker会停止,但在高负载、短超时场景下,Goroutine数量持续攀升且无法恢复。相关代码如下:
func processJobs(ctx context.Context, jobs []Job) ([]Result, error) { results := make(chan Result) // unbuffered var wg sync.WaitGroup for _, job := range jobs { wg.Add(1) go func(j Job) { defer wg.Done() select { case results <- doWork(j): case <-ctx.Done(): return } }(job) } go func() { wg.Wait() close(results) }() var out []Result for { select { case r, ok := <-results: if !ok { return out, nil } out = append(out, r) case <-ctx.Done(): return out, ctx.Err() } } }
泄漏成因
- 未缓冲通道的发送阻塞:当上下文取消时,主函数会立刻返回,不再接收
results通道的数据。但部分worker可能已经完成了doWork(j)的执行,正尝试往未缓冲的results通道发送结果——未缓冲通道的发送操作必须等待接收方才能完成,此时没有接收方,这些worker会永久阻塞在发送操作上,无法退出,造成Goroutine泄漏。 - select分支的触发时机问题:worker中的
select只有在doWork(j)执行完成前上下文取消,才会走到ctx.Done()分支退出。如果doWork(j)已经执行完毕,results <- doWork(j)分支会处于就绪状态,select会优先选择这个分支,此时即使上下文已取消,worker也会卡在发送操作上,无法响应ctx.Done()。
修复方案
核心解决思路是避免worker在上下文取消后阻塞在结果发送操作上,以下是可靠的修复实现:
func processJobs(ctx context.Context, jobs []Job) ([]Result, error) { // 给结果通道添加与任务数一致的缓冲,避免worker发送结果阻塞 results := make(chan Result, len(jobs)) var wg sync.WaitGroup for _, job := range jobs { wg.Add(1) go func(j Job) { defer wg.Done() // 先检查上下文是否已取消,避免无效执行任务 select { case <-ctx.Done(): return default: } res := doWork(j) // 任务完成后,再次检查上下文,确保发送结果时不会阻塞 select { case results <- res: case <-ctx.Done(): return } }(job) } go func() { wg.Wait() close(results) }() var out []Result for { select { case r, ok := <-results: if !ok { return out, nil } out = append(out, r) case <-ctx.Done(): // 上下文取消后,启动后台goroutine消费剩余结果,确保worker能正常退出 go func() { for range results { } }() return out, ctx.Err() } } }
关键修复点说明
- 带缓冲的结果通道:缓冲大小等于任务数,确保worker发送结果时不会因为没有接收方而阻塞,即使主函数提前退出,缓冲也能容纳所有已完成任务的结果,worker可以顺利完成发送后退出。
- 双层select检查:worker先检查上下文再执行任务,避免无效计算;任务完成后再次检查上下文,确保发送结果时如果上下文已取消,直接退出而不阻塞。
- 后台消费剩余结果:主函数在上下文取消时,启动goroutine消费
results通道中剩余的结果,彻底避免worker因发送结果而阻塞。
内容的提问来源于stack exchange,提问作者Md.Jewel Mia
相关产品推荐
相关产品推荐

