Go如何地道实现goroutine并发执行与全量结果/错误收集
问题描述
在下方代码块中,我尝试启动多个goroutine执行任务,并获取所有协程的返回结果(包含成功结果与错误信息):
package main import ( "fmt" "sync" ) func processBatch(num int, errChan chan<- error, resultChan chan<- int, wg *sync.WaitGroup) { defer wg.Done() if num == 3 { resultChan <- 0 errChan <- fmt.Errorf("goroutine %d's error returned", num) } else { square := num * num resultChan <- square errChan <- nil } } func main() { var wg sync.WaitGroup batches := [5]int{1, 2, 3, 4, 5} resultChan := make(chan int) errChan := make(chan error) for i := range batches { wg.Add(1) go processBatch(batches[i], errChan, resultChan, &wg) } var results [5]int var err [5]error for i := range batches { results[i] = <-resultChan err[i] = <-errChan } wg.Wait() close(resultChan) close(errChan) fmt.Println(results) fmt.Println(err) }
代码可在Go Playground上运行。
上述代码可以正常运行,得到符合预期的输出如下:
[25 1 4 0 16] [<nil> <nil> <nil> goroutine 3's error returned <nil>]
我想了解是否存在更符合Go语言惯用法的实现方式来达成该需求。我查阅了errgroup包的官方文档,但未找到可适配该场景的用法,欢迎提供相关实现建议。
符合Go惯用法的实现方案
原有写法可以运行,但存在三个可优化点:
- 拆分
resultChan和errChan两个无缓冲通道,需要严格保证读写顺序一一匹配,一旦逻辑调整很容易触发死锁 - 任务函数
processBatch耦合了WaitGroup、通道等并发原语,函数本身的可测试性差 wg.Wait()调用时机晚于通道读取,若协程异常未写入通道会触发永久阻塞
没找到errgroup适配方案是正常的:标准errgroup的设计目标是一组任务任意一个失败就取消整体执行,天生不适合"跑完所有任务、收集每个任务成功/失败结果"的场景,不需要硬套。
根据场景不同,有两种非常地道的实现方式:
方案1:固定任务数场景(最简单、性能最高)
如果提前明确任务总数量,不需要流式处理结果,可以直接预分配结果/错误切片,每个协程只写入自己对应的索引位,完全不需要通道,不存在并发写冲突:
package main import ( "fmt" "sync" ) // 任务逻辑完全剥离并发原语,是纯粹的输入输出函数,可直接单测 func processBatch(num int) (int, error) { if num == 3 { return 0, fmt.Errorf("goroutine %d's error returned", num) } return num * num, nil } func main() { batches := [5]int{1, 2, 3, 4, 5} var ( wg sync.WaitGroup results [5]int errs [5]error ) for i, num := range batches { wg.Add(1) // 把索引和参数通过参数传入协程,避免闭包变量捕获问题 go func(idx int, n int) { defer wg.Done() results[idx], errs[idx] = processBatch(n) }(i, num) } wg.Wait() fmt.Println(results) fmt.Println(errs) }
这个写法的优势:
- 无通道额外开销,性能最高
- 结果顺序和输入顺序完全一致,不会因为协程调度先后乱序
- 代码逻辑简洁,出错概率极低
方案2:通用流式场景(适配动态任务数、结果流式处理)
如果任务数量动态变化,或者需要边执行边处理结果,可以用单通道打包返回值的方式,避免多通道匹配问题:
package main import ( "fmt" "sync" ) // TaskResult 打包单个任务的输入、输出、错误,和通道绑定 type TaskResult struct { Input int Output int Err error } func processBatch(num int) (int, error) { if num == 3 { return 0, fmt.Errorf("goroutine %d's error returned", num) } return num * num, nil } func main() { batches := [5]int{1, 2, 3, 4, 5} // 初始化和任务数等长的缓冲通道,协程写入不会阻塞 resultChan := make(chan TaskResult, len(batches)) var wg sync.WaitGroup for _, num := range batches { wg.Add(1) go func(n int) { defer wg.Done() res, err := processBatch(n) resultChan <- TaskResult{ Input: n, Output: res, Err: err, } }(num) } // 单独起协程等待所有任务完成后关闭通道,避免遍历通道时永久阻塞 go func() { wg.Wait() close(resultChan) }() // 直接遍历通道拿所有结果,不需要手动控制读取次数 var results [5]int var errs [5]error for r := range resultChan { // 按输入值映射索引,保证结果顺序和输入一致 results[r.Input-1] = r.Output errs[r.Input-1] = r.Err } fmt.Println(results) fmt.Println(errs) }
这个写法的优势:
- 单通道传输不存在多通道收发顺序错配的死锁风险
- 支持动态增减任务数量,不需要修改结果读取逻辑
- 可以在遍历通道时实时处理已完成的任务结果,不需要等所有任务跑完
内容的提问来源于stack exchange,提问作者Anshul Goyal
相关产品推荐
相关产品推荐

