使用errgroup实现Go WorkerPool时协程阻塞,g.Wait()未触发
问题分析与解决
阻塞根源
你的代码存在两个核心问题,导致错误发生时主goroutine直接阻塞,永远执行不到g.Wait():
- 主goroutine硬编码要从
results读取totalUsers次数据,但当某个worker因CallUserAPI报错退出时,对应任务的结果根本不会写入results,主goroutine会一直卡在r := <-results这一步,永远等不到足够的结果。 workerUser处理jobs通道时,没有判断通道是否关闭——当jobs被关闭后,user, _ := <-jobs会持续返回零值,worker会无意义地循环调用CallUserAPI,直到ctx被取消。
修复方案
1. 修正Worker函数
处理jobs通道的关闭事件,同时往results写数据前先检查ctx状态,避免写入阻塞:
func workerUser(ctx context.Context, jobs <-chan UsersInfo, results chan<- UsersInfo) error { for { select { case <-ctx.Done(): return ctx.Err() case user, ok := <-jobs: // 通道关闭时直接退出循环 if !ok { return nil } userInfo, err := CallUserAPI(ctx, user) if err != nil { return err } // 写入结果前先确认ctx未被取消 select { case <-ctx.Done(): return ctx.Err() case results <- userInfo: } } } }
2. 修正主goroutine逻辑
不再固定读取totalUsers次结果,改为同时监听results的完成信号和errgroup的错误信号,实现错误发生时立即停止等待:
jobs := make(chan UsersInfo, totalUsers) results := make(chan UsersInfo, totalUsers) g, gCtx := errgroup.WithContext(ctx) // 启动worker for i := 1; i <= 4; i++ { g.Go(func() error { return workerUser(gCtx, jobs, results) }) } // 单独开goroutine写入jobs,避免主goroutine阻塞在写入操作 go func() { for _, user := range usersResp { select { case <-gCtx.Done(): return case jobs <- user: } } close(jobs) }() var usersArray []UsersInfo // 启动goroutine接收所有结果 done := make(chan struct{}) go func() { for r := range results { usersArray = append(usersArray, r) } close(done) }() // 等待errgroup完成或所有结果接收完毕 select { case <-done: case <-gCtx.Done(): } // 关闭results通道,避免接收goroutine泄漏 close(results) // 最后统一处理错误 if err := g.Wait(); err != nil { return nil, err }
关键优化点
- 主goroutine通过双信号监听,不再硬等固定数量的结果,错误发生时能立即响应。
- worker处理jobs通道的关闭事件,避免无效循环。
- 写入results前检查ctx状态,防止worker因主goroutine阻塞而卡在写入步骤。
内容的提问来源于stack exchange,提问作者ARKhan
相关产品推荐
相关产品推荐

