如何在所有工作完成后关闭多goroutine写入的channel?
解决Go中任务结果通道的关闭问题
要解决最后一个循环阻塞的问题,核心是等所有worker都处理完任务并返回结果后,再关闭producedResults通道。这里可以用sync.WaitGroup来跟踪所有worker的执行状态,具体实现步骤如下:
- 引入
sync包,创建WaitGroup实例,用来统计活跃的worker数量 - 启动每个worker前,调用
wg.Add(1)增加计数 - 每个worker在处理完所有任务(退出
for taskSize := range remainingTasks循环后),调用wg.Done()减少计数 - 单独启动一个goroutine,等待所有worker完成(
wg.Wait())后,关闭producedResults通道
修改后的完整代码:
package main import ( "sync" "time" ) func main() { quantityOfTasks := 100 quantityOfWorkers := 60 remainingTasks := make(chan int) producedResults := make(chan int) var wg sync.WaitGroup // 生产任务 go func() { for i := 0; i < quantityOfTasks; i++ { remainingTasks <- 1 } close(remainingTasks) }() // 启动worker for i := 0; i < quantityOfWorkers; i++ { wg.Add(1) // 每个worker启动前增加计数 go func() { defer wg.Done() // worker退出时减少计数,用defer确保执行时机 for taskSize := range remainingTasks { // 模拟耗时任务 time.Sleep(time.Second * time.Duration(taskSize)) // 返回结果 producedResults <- taskSize } }() } // 等待所有worker完成后关闭结果通道 go func() { wg.Wait() close(producedResults) }() // 汇总结果 executedTasks := 0 for resultOfTheTask := range producedResults { executedTasks += resultOfTheTask } }
关键说明:
sync.WaitGroup的作用是等待所有worker goroutine执行完毕,确保不会提前关闭producedResults导致结果丢失defer wg.Done()放在worker goroutine开头,能保证不管worker是正常退出还是遇到异常,都会正确减少计数- 单独的goroutine负责等待和关闭通道,不会阻塞主goroutine的结果汇总逻辑
内容的提问来源于stack exchange,提问作者AntonioR
相关产品推荐
相关产品推荐

