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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.30 04:33:25