【问题标题】:Get responses from multiple go routines into an array从多个 go 例程中获取响应到一个数组中
【发布时间】:2019-06-29 06:46:47
【问题描述】:

我需要从多个 go 例程中获取响应并将它们放入一个数组中。我知道可以为此使用通道,但是我不确定如何确保所有 go 例程都已完成对结果的处理。因此我正在使用等待组。

代码

func main() {
  log.Info("Collecting ints")
  var results []int32
  for _, broker := range e.BrokersByBrokerID {
      wg.Add(1)
      go getInt32(&wg)
  }
  wg.Wait()
  log.info("Collected")
}

func getInt32(wg *sync.WaitGroup) (int32, error) {
  defer wg.Done()

  // Just to show that this method may just return an error and no int32
  err := broker.Open(config)
  if err != nil && err != sarama.ErrAlreadyConnected {
    return 0, fmt.Errorf("Cannot connect to broker '%v': %s", broker.ID(), err)
  }
  defer broker.Close()

  return 1003, nil
}

我的问题

如何将所有响应 int32(可能返回错误)放入我的 int32 数组,确保所有 go 例程都已完成处理工作并返回错误或 int?

【问题讨论】:

  • 渠道是正确的答案。

标签: go concurrency goroutine


【解决方案1】:

如果您不处理作为 goroutine 启动的函数的返回值,它们将被丢弃。见What happens to return value from goroutine

您可以使用切片来收集结果,其中每个 goroutine 可以接收将结果放入的索引,或者元素的地址。见Can I concurrently write different slice elements。请注意,如果您使用它,则必须预先分配切片,并且只能写入属于 goroutine 的元素,您不能“触摸”其他元素,也不能附加到切片。

或者您可以使用一个通道,goroutine 在该通道上发送包含它们处理的项目的索引或 ID 的值,因此收集的 goroutine 可以识别或排序它们。见How to collect values from N goroutines executed in a specific order?

如果处理应在遇到第一个错误时停止,请参阅Close multiple goroutine if an error occurs in one in go

以下是使用频道时的示例。请注意,这里不需要等待组,因为我们知道我们希望通道上的值与我们启动的 goroutine 一样多。

type result struct {
    task int32
    data int32
    err  error
}

func main() {
    tasks := []int32{1, 2, 3, 4}

    ch := make(chan result)

    for _, task := range tasks {
        go calcTask(task, ch)
    }

    // Collect results:
    results := make([]result, len(tasks))

    for i := range results {
        results[i] = <-ch
    }

    fmt.Printf("Results: %+v\n", results)
}

func calcTask(task int32, ch chan<- result) {
    if task > 2 {
        // Simulate failure
        ch <- result{task: task, err: fmt.Errorf("task %v failed", task)}
        return
    }

    // Simulate success
    ch <- result{task: task, data: task * 2, err: nil}
}

输出(试试Go Playground):

Results: [{task:4 data:0 err:0x40e130} {task:1 data:2 err:<nil>} {task:2 data:4 err:<nil>} {task:3 data:0 err:0x40e138}]

【讨论】:

  • 需要注意的是,如果使用slice,不小心将append 转为slice 会导致数据争用。因此,在这种情况下,我会推荐array
  • @leafbebop 是的,如果使用切片,它必须预先分配,并且只有给定索引处的元素可以写入 goroutine。
  • 您介意编辑我给定的代码示例,使其适用于您建议的渠道方法吗?由于等待组,我有点困惑。我相信拥有代码解决方案将是最容易理解的,因为我不确定我是否应该在您提出的解决方案中使用等待组
  • @kentor Waitgroup 可用于等待一组 goroutine 完成工作。如果你使用一个频道并且你知道你期望它有多少值,你甚至不需要 WaitGroup,你可以从频道中读取n 值,然后你就知道你已经完成了。
  • @kentor “失败”的 Goroutines 也应该返回一个值,所以是的,你期望与你启动的 goroutines 一样多的值。在通道上发送的值应该是结构值,包装当前函数的 2 个返回值。所以即使应该返回一个错误,那仍然是一个要在通道上发送的值。
【解决方案2】:

我也相信你必须使用频道,它必须是这样的:

package main

import (
    "fmt"
    "log"
    "sync"
)

var (
    BrokersByBrokerID = []int32{1, 2, 3}
)

type result struct {
    data string
    err string // you must use error type here
}

func main()  {
    var wg sync.WaitGroup
    var results []result
    ch := make(chan result)

    for _, broker := range BrokersByBrokerID {
        wg.Add(1)
        go getInt32(ch, &wg, broker)
    }

    go func() {
        for v := range ch {
            results = append(results, v)
        }
    }()

    wg.Wait()
    close(ch)

    log.Printf("collected %v", results)
}

func getInt32(ch chan result, wg *sync.WaitGroup, broker int32) {
    defer wg.Done()

    if broker == 1 {
        ch <- result{err: fmt.Sprintf("error: gor broker 1")}
        return
    }

    ch <- result{data: fmt.Sprintf("broker %d - ok", broker)}
}

结果将如下所示:

2019/02/05 15:26:28 collected [{broker 3 - ok } {broker 2 - ok } { error: gor broker 1}]

【讨论】:

  • 我没有考虑创建一个 goroutine 来消费结果。感谢那!我假设不需要互斥锁,因为只有一个 go 例程写入数组?
  • 没有更高层次的方法可以做到这一点吗?
猜你喜欢
  • 1970-01-01
  • 2020-04-08
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2019-05-08
  • 2021-09-14
  • 2022-11-27
相关资源
最近更新 更多