You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何在所有工作完成后关闭多goroutine写入的channel?

解决Go中任务结果通道的关闭问题

要解决最后一个循环阻塞的问题,核心是等所有worker都处理完任务并返回结果后,再关闭producedResults通道。这里可以用sync.WaitGroup来跟踪所有worker的执行状态,具体实现步骤如下:

  1. 引入sync包,创建WaitGroup实例,用来统计活跃的worker数量
  2. 启动每个worker前,调用wg.Add(1)增加计数
  3. 每个worker在处理完所有任务(退出for taskSize := range remainingTasks循环后),调用wg.Done()减少计数
  4. 单独启动一个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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.11 19:19:59