使用sync.WaitGroup与缓冲Channel的并发同步问题排查
Go中sync.WaitGroup与缓冲Channel配合的问题及修复
问题描述
使用sync.WaitGroup配合缓冲Channel时遇到两个核心问题:
- 主goroutine在
WaitGroup完成后立即关闭Channel,导致读取goroutine尚未读完Channel内所有数据,最终统计的学生数据不完整。 - 程序遇到错误时无法立即终止所有执行流程,仍会继续发起不必要的HTTP请求。
问题分析
- 数据读取中断原因:
wg.Wait()仅等待所有数据写入goroutine完成,但主goroutine在等待结束后立刻关闭studentCh并统计students切片长度——此时读取goroutine可能还在从Channel中读取数据并追加到切片,导致最终统计的长度小于实际获取到的数据量。 - 错误无法立即终止原因:当前读取goroutine收到错误后仅跳出循环,未通知其他goroutine停止执行,也未终止整个程序,导致后续仍会有无效请求发起。
修复方案
核心修改点
- 新增
readWg等待读取goroutine完成所有数据处理,确保统计结果准确。 - 使用
context.WithCancel创建可取消上下文,遇到错误时取消上下文,终止所有递归goroutine的执行。 - 优化错误处理逻辑:收到错误后立即取消上下文、关闭Channel,并终止程序流程。
修改后的代码
主函数
func main() { var wg sync.WaitGroup var readWg sync.WaitGroup // 新增:等待读取goroutine完成数据处理 start := time.Now() students := make([]studentDetail, 0) studentCh := make(chan studentDetail, 10000) errorCh := make(chan error, 1) // 创建可取消上下文,用于终止所有goroutine rCtx, cancel := context.WithCancel(context.Background()) defer cancel() wg.Add(1) go s.getDetailStudents(rCtx, studentCh, errorCh, &wg, s.Link, false) readWg.Add(1) go func(ch chan studentDetail, e chan error) { defer readWg.Done() LOOP: for { select { case p, ok := <-ch: if ok { L.Printf("Links %s: [%s]\n", p.title, p.link) students = append(students, p) } else { L.Print("Closed data channel") break LOOP } case err := <-e: if err != nil { L.Errorf("Encountered error: %v", err) cancel() // 取消上下文,终止所有goroutine close(ch) // 提前关闭数据Channel break LOOP } case <-rCtx.Done(): // 上下文已取消,直接退出循环 break LOOP } } }(studentCh, errorCh) // 等待所有数据写入goroutine完成 wg.Wait() close(studentCh) // 所有写入操作完成后关闭Channel close(errorCh) // 等待读取goroutine处理完所有剩余数据 readWg.Wait() L.Warnln("All operations completed!") L.Warnf("total items fetched: %d", len(students)) elapsed := time.Since(start) L.Warnf("operation took %s", elapsed) }
递归函数
func (s Student) getDetailStudents(rCtx context.Context, content chan<- studentDetail, errorCh chan<- error, wg *sync.WaitGroup, url string, subSection bool) { util.MustNotNil(rCtx) L := logger.GetLogger(rCtx) defer func() { L.Println("WaitGroup done for url:", url) wg.Done() }() // 先检查上下文是否已取消,避免发起无效请求 select { case <-rCtx.Done(): L.Warnf("Context canceled, skipping request for %s", url) return default: } wc := getWC() httpClient := wc.Registry.MustHTTPClient() res, err := httpClient.Get(url) if err != nil { L.Errorf("Request failed for %s: %v", url, err) // 非阻塞发送错误,避免goroutine阻塞 select { case errorCh <- err: default: } return } defer res.Body.Close() if res.StatusCode != 200 { errMsg := fmt.Sprintf("status code error: %d %s for url %s", res.StatusCode, res.Status, url) L.Error(errMsg) select { case errorCh <- errors.New(errMsg): default: } return } // 判断是否需要发起递归请求 if !subSection { wg.Add(1) go s.getDetailStudents(rCtx, content, errorCh, wg, link, true) } // 检查上下文状态,避免无效解析操作 select { case <-rCtx.Done(): L.Warnf("Context canceled, skipping parsing for %s", url) return default: } // 解析学生数据并发送到Channel students := s.parseStudentItemList(rCtx, item) for _, student := range students { select { case content <- student: case <-rCtx.Done(): // 上下文取消,停止发送数据 return } } L.Warnf("Fetched %d records from %q", len(students), url) }
关键说明
- 读取goroutine等待:新增
readWg确保主goroutine在读取goroutine处理完所有Channel数据后再统计结果,彻底避免数据丢失。 - 上下文取消机制:通过
context.WithCancel实现全局终止信号,所有递归goroutine会定期检查上下文状态,及时终止无效操作。 - 非阻塞错误发送:使用
select发送错误到errorCh,避免因Channel已满导致goroutine阻塞。
内容的提问来源于stack exchange,提问作者Sankalp
相关产品推荐
相关产品推荐

