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

为何Goroutine并发任务仅返回部分错误而非全部错误?

并发任务错误合并丢失问题

我定义了Cycle类用于处理并发任务,期望将Do1、Do2两个函数分别放入Goroutine中运行,等待全部任务完成后合并所有返回的错误,但最终仅能得到ERR1或ERR2其中一个错误,无法获取两者的组合结果。各方法职责如下:

  • Run:在Goroutine中运行函数并收集其错误
  • WaitAllDone:合并所有函数的错误并等待全部任务完成
  • Do1、Do2:测试函数
import (
    "fmt"
    "go.uber.org/multierr"
    "sync"
    "testing"
)

type Cycle struct {
    errChan chan error
    wg sync.WaitGroup
}

func NewCycle() *Cycle {
    return &Cycle{
        errChan: make(chan error),
        wg:      sync.WaitGroup{},
    }
}

// run fn and collect its error into error channel
func (c *Cycle) Run(fn func() error) {
    c.wg.Add(1)
    go func() {
        defer c.wg.Done()
        if err := fn(); err != nil {
            c.errChan <- err
        }
    }()
}

// wait all fn finish and combine their error together
func (c *Cycle) WaitAllDone() error {
    var err error
    go func() {
        for {
            if tmpErr, ok := <-c.errChan; ok {
                err = multierr.Append(err, tmpErr)
            } else{
                break
            }
        }
    }()
    c.wg.Wait()
    close(c.errChan)
    return err
}

func Do1() error {
    return fmt.Errorf("ERR1")
}

func Do2() error {
    return fmt.Errorf("ERR2")
}

func Test41(t *testing.T) {
    c := NewCycle()
    c.Run(Do1)
    c.Run(Do2)
    if err := c.WaitAllDone(); err != nil {
        t.Log(err)
    }
}

最终t.Log(err)仅输出ERR1或ERR2,而我期望输出ERR1 ERR2,请问为何会丢失一个错误?


问题原因

WaitAllDone方法存在竞态条件:

  1. 主goroutine调用c.wg.Wait()等待所有任务完成后,立即关闭errChan并返回err变量。
  2. 用于读取错误并合并的后台goroutine,并没有被主goroutine等待。当主goroutine返回时,后台goroutine可能还没来得及读取第二个错误并完成合并,导致最终返回的错误只包含其中一个。

修复方案

方案1:等待错误合并完成

给读取错误的goroutine添加等待组,确保所有错误都被合并后再返回:

func (c *Cycle) WaitAllDone() error {
    var err error
    var readWg sync.WaitGroup
    readWg.Add(1)
    go func() {
        defer readWg.Done()
        for tmpErr := range c.errChan {
            err = multierr.Append(err, tmpErr)
        }
    }()
    c.wg.Wait()
    close(c.errChan)
    readWg.Wait() // 等待错误读取合并完成
    return err
}

方案2:用互斥锁直接保护错误变量

放弃channel,改用互斥锁保证错误合并的线程安全,逻辑更简洁:

type Cycle struct {
    mu  sync.Mutex
    err error
    wg  sync.WaitGroup
}

func NewCycle() *Cycle {
    return &Cycle{}
}

func (c *Cycle) Run(fn func() error) {
    c.wg.Add(1)
    go func() {
        defer c.wg.Done()
        if tmpErr := fn(); tmpErr != nil {
            c.mu.Lock()
            c.err = multierr.Append(c.err, tmpErr)
            c.mu.Unlock()
        }
    }()
}

func (c *Cycle) WaitAllDone() error {
    c.wg.Wait()
    return c.err
}

两种方案都能确保所有错误被正确合并,最终输出ERR1 ERR2。

内容的提问来源于stack exchange,提问作者hankeyyh

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.20 15:32:02