如何在任意数量goroutine完成时实时处理其输出结果
嵌套goroutine场景下实时处理输出的实现方案
核心问题本质是动态生成goroutine时,计数操作存在并发安全问题、且计数时机不对导致总任务数统计失准,解决方案可以直接用sync.WaitGroup配合结果通道实现,完全满足实时处理输出+等待所有任务完成的需求:
实现代码
package main import "sync" func main() { // 初始化WaitGroup和结果传输通道 var wg sync.WaitGroup resChan := make(chan string) initUrls := []string{"https://example.com", "..."} // 替换为你的初始URL列表 // 第一层goroutine启动前先加计数,Add必须在启动goroutine前调用 wg.Add(len(initUrls)) for _, u := range initUrls { go func(url string) { defer wg.Done() // 当前goroutine执行完成后减计数 // 第一层请求逻辑 data := get(url) // 第一层结果写入通道,可实时被接收处理 resChan <- data // 解析得到子URL列表 subUrls := getUrls(data) // 启动子goroutine前,先把对应数量加到WaitGroup计数中 wg.Add(len(subUrls)) for _, su := range subUrls { go func(subUrl string) { defer wg.Done() subData := get(subUrl) // 子goroutine结果也写入通道 resChan <- subData }(su) } }(u) } // 单独启动一个goroutine等待所有任务执行完成后关闭结果通道 go func() { wg.Wait() close(resChan) }() // 遍历结果通道即可实时处理所有返回的结果,通道关闭后循环自动退出 for res := range resChan { // 这里写你处理结果的逻辑,拿到一个处理一个 processResult(res) } } // 以下为示例用到的模拟方法,替换为你自己的实现即可 func get(url string) string { // 请求URL返回数据的逻辑 return "" } func getUrls(data string) []string { // 从返回数据中解析子URL的逻辑 return nil } func processResult(res string) { // 实时处理结果的逻辑 }
关键要点
sync.WaitGroup的Add()方法必须在启动对应goroutine的父逻辑中调用,绝对不能放到goroutine内部执行,否则可能出现主流程已经执行到Wait()时,子goroutine还没调用Add(),导致计数遗漏提前结束等待- 单独开goroutine等待所有任务完成后关闭结果通道,主流程不需要提前知道总任务数,遍历通道即可实现实时处理,也不会出现通道永远阻塞的问题
- 该方案支持任意深度的goroutine嵌套,只要每启动下一级goroutine之前调用对应数量的
Add()即可正常工作
原方案失效原因说明
你之前把计数rc++放在goroutine内部的写法存在两个问题:
- 普通整型变量
rc并发读写没有加同步机制,存在竞态条件,计数结果不准确 - 主流程读取
rc计算总任务数时,部分goroutine可能还没执行到rc++的逻辑,拿到的总任务数比实际小,会导致主流程提前退出,剩余子goroutine写入通道时永远阻塞,出现goroutine泄漏
内容的提问来源于stack exchange,提问作者beardeadclown
相关产品推荐
相关产品推荐

