在sync.WaitGroup协程中写入chan:为何需单独协程关闭通道?
为什么必须用单独协程执行
wg.Wait()后关闭通道? 我从API获取条目列表,再为每个条目发起单独API请求拿详情。因为API令牌有速率限制,没法并发请求每个条目,但初始API响应支持分页,所以可以并发处理分页数据。我找到两种代码实现,一种能正常运行,另一种不行,原因如下:
可行代码
func main() { // 模拟初始API返回的分页结果 page1 := []int{1, 2, 3} page2 := []int{4, 5, 6} page3 := []int{7, 8, 9} pages := [][]int{page1, page2, page3} results := make(chan string) var wg sync.WaitGroup for i := range pages { wg.Add(1) go func(i int) { defer wg.Done() for j := range pages[i] { // 模拟发起额外API请求并生成报告 time.Sleep(500 * time.Millisecond) result := fmt.Sprintf("Finished creating report for %d", pages[i][j]) results <- result } }(i) } // 单独协程等待所有任务完成后关闭通道 go func() { wg.Wait() close(results) }() // 持续接收通道结果直到通道关闭 for result := range results { fmt.Println(result) } }
无效代码
func main() { // 模拟初始API返回的分页结果 page1 := []int{1, 2, 3} page2 := []int{4, 5, 6} page3 := []int{7, 8, 9} pages := [][]int{page1, page2, page3} results := make(chan string) var wg sync.WaitGroup for i := range pages { wg.Add(1) go func(i int) { defer wg.Done() for j := range pages[i] { // 模拟发起额外API请求并生成报告 time.Sleep(500 * time.Millisecond) result := fmt.Sprintf("Finished creating report for %d", pages[i][j]) results <- result } }(i) } // 此处写法会导致死锁 wg.Wait() close(results) // 永远执行不到这里 for result := range results { fmt.Println(result) } }
核心原因:无缓冲通道的发送阻塞问题
Go的无缓冲通道(make(chan string)默认是无缓冲)的发送操作会阻塞,直到有接收者接收数据。
无效代码的死锁逻辑:
- 启动所有分页协程后,main函数立刻执行
wg.Wait(),等待所有协程完成 - 分页协程执行到
results <- result时,因为main还在wg.Wait(),没有进入for range接收数据,发送操作被阻塞 - 所有分页协程都卡在发送步骤,永远无法执行到
defer wg.Done(),wg.Wait()永远等不到所有任务完成,程序陷入死锁
可行代码的正确逻辑:
- 启动分页协程的同时,启动一个独立协程负责等待所有任务完成后关闭通道
- main函数立刻进入
for range results开始接收数据,通道有了接收者,分页协程的发送操作不会阻塞,能正常执行到wg.Done() - 当所有分页协程完成,独立协程的
wg.Wait()返回,关闭results通道 - main的
for range接收到通道关闭的信号,循环结束,程序正常退出
内容的提问来源于stack exchange,提问作者ETLJ
相关产品推荐
相关产品推荐

