【问题标题】:How to use channels to gather response from various goroutines如何使用通道收集来自各种 goroutine 的响应
【发布时间】:2020-09-01 09:47:50
【问题描述】:

我是 Golang 的新手,我有一个使用 WaitGroupMutex 实现的任务,我想将其转换为使用 Channels

对该任务的一个非常简短的描述是:根据需要摒弃尽可能多的 go 例程来处理结果,并在主 go 例程中等待并收集所有结果。

我使用WaitGroupMutex的实现如下:

package main

import (
    "fmt"
    "math/rand"
    "sync"
    "time"
)

func process(input int, wg *sync.WaitGroup, result *[]int, lock *sync.Mutex) *[]int {
    defer wg.Done()
    defer lock.Unlock()

    rand.Seed(time.Now().UnixNano())
    n := rand.Intn(5)
    time.Sleep(time.Duration(n) * time.Second)
    lock.Lock()
    *result = append(*result, input * 10)

    return result
}

func main() {

    var wg sync.WaitGroup
    var result []int
    var lock sync.Mutex
    for i := range []int{1,2,3,4,5} {
        wg.Add(1)
        go process(i, &wg, &result, &lock)
    }
}

如何将使用Mutex 的内存同步替换为使用Channels 的内存同步?

我的主要问题是我不确定如何确定处理最终任务的最终 go 例程,因此让那个例程成为关闭 channel 的例程。这个想法是,通过关闭 channel,主 go 例程可以循环遍历 channel,检索结果,当它看到 channel 已关闭时,它会继续前进。

在这种情况下,关闭频道的方法也可能是错误的,因此我在这里问。

更有经验的 Go 程序员如何使用 channels 解决这个问题?

【问题讨论】:

  • 关闭通道是正确的。因为您有多个发送者,所以您仍然需要 WaitGroup,以及一个在 WaitGroup 完成后关闭通道的额外 goroutine。
  • 创建一个通道并将其传递给process,以便它可以发送结果。使用另一个例程等待WaitGroup,然后关闭通道。同时,主要使用range从通道读取结果。
  • @hmm 介意在答案中用代码勾勒出来吗?谢谢
  • 与您的问题无关,但有一些建议。在许多 goroutine 中使用 rand 包的预制 rng 会导致对 rng 的争用。为了获得更高的性能(如果您的应用程序需要),请在每个 goroutine 中使用 rand.New(rand.NewSource( yourSeed )) 实例化一个新的 rng。此外,无需在每个 goroutine 中“重新播种”预制的 rng ——只需在程序开始时播种一次即可。

标签: go concurrency coroutine goroutine


【解决方案1】:

我更改了您的代码以使用该频道。还有很多其他方法可以使用该频道。

package main

import (
    "fmt"
    "math/rand"
    "time"
)

func process(input int, out chan<- int) {
    rand.Seed(time.Now().UnixNano())
    n := rand.Intn(5)
    time.Sleep(time.Duration(n) * time.Second)
    out <- input * 10
}

func main() {
    var result []int
    resultChan := make(chan int)
    items := []int{1, 2, 3, 4, 5}

    for _, v := range items {
        go process(v, resultChan)
    }

    for i := 0; i < len(items); i++ {
        res, _ := <-resultChan
        result = append(result, res)
    }

    close(resultChan)
    fmt.Println(result)
}

更新:(评论的答案)

如果项目数未知,您需要向主发出信号以完成。否则“死锁”,您可以创建一个通道来指示主要功能完成。也可以使用sync.waiteGroup

对于 Goroutine 中的 panic,你可以使用 defer 和 recovery 来处理错误。并且您可以创建一个错误通道矿石,您可以使用x/sync/errgroup

有很多解决方案。这取决于你的问题。所以没有具体的方式来使用 goroutine、channel 和...

【讨论】:

  • 这个解决方案利用了迭代次数已知的知识。即通过使用 len(items) 如果这个数字未知,将如何处理?
  • 如果其中一个 go 例程出现恐慌会发生什么?它不会通过循环吗?
【解决方案2】:

这是一个示例 sn-p,其中我使用通道切片而不是等待组来执行分叉连接:

package main

import (
    "fmt"
    "os"
)

type cStruct struct {
    resultChan chan int
    errChan    chan error
}

func process(i int) (v int, err error) {
    v = i
    return
}

func spawn(i int) cStruct {
    r := make(chan int)
    e := make(chan error)
    go func(i int) {
        defer close(r)
        defer close(e)
        v, err := process(i)
        if err != nil {
            e <- err
            return
        }
        r <- v
        return
    }(i)
    return cStruct{
        r,
        e,
    }
}

func main() {
    //have a slice of channelStruct
    var cStructs []cStruct
    nums := []int{1, 2, 3, 4, 5}
    for _, v := range nums {
        cStruct := spawn(v)
        cStructs = append(cStructs, cStruct)
    }
    //All the routines have been spawned, now iterate over the slice:
    var results []int
    for _, c := range cStructs {
        rChan, errChan := c.resultChan, c.errChan
        select {
        case r := <-rChan:
            {
                results = append(results, r)
            }
        case err := <-errChan:
            {
                if err != nil {
                    os.Exit(1)
                    return
                }
            }
        }

    }
    //All the work should be done by now, iterating over the results
    for _, result := range results {
        fmt.Println("Aggregated result:", result)
    }
}

【讨论】:

    【解决方案3】:

    这是一个使用WaitGroup 的解决方案,而不是等待固定数量的结果。

    package main
    
    import (
        "fmt"
        "math/rand"
        "sync"
        "time"
    )
    
    func process(input int, wg *sync.WaitGroup, resultChan chan<- int) {
        defer wg.Done()
    
        rand.Seed(time.Now().UnixNano())
        n := rand.Intn(5)
        time.Sleep(time.Duration(n) * time.Second)
    
        resultChan <- input * 10
    }
    
    func main() {
        var wg sync.WaitGroup
    
        resultChan := make(chan int)
    
        for i := range []int{1,2,3,4,5} {
            wg.Add(1)
            go process(i, &wg, resultChan)
        }
    
        go func() {
            wg.Wait()
            close(resultChan)
        }()
    
        var result []int
        for r := range resultChan {
            result = append(result, r)
        }
    
        fmt.Println(result)
    }
    

    【讨论】:

      猜你喜欢
      • 2021-06-26
      • 2017-03-26
      • 2013-04-28
      • 1970-01-01
      • 2016-03-16
      • 2016-10-31
      • 1970-01-01
      • 2023-04-02
      • 1970-01-01
      相关资源
      最近更新 更多