You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

扇出/扇入工作池上下文取消后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()
        }
    }
}
泄漏成因
  1. 未缓冲通道的发送阻塞:当上下文取消时,主函数会立刻返回,不再接收results通道的数据。但部分worker可能已经完成了doWork(j)的执行,正尝试往未缓冲的results通道发送结果——未缓冲通道的发送操作必须等待接收方才能完成,此时没有接收方,这些worker会永久阻塞在发送操作上,无法退出,造成Goroutine泄漏。
  2. 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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.16 09:33:05