【问题标题】:How to perform concurrent downloads in Go如何在 Go 中执行并发下载
【发布时间】:2015-12-02 23:06:41
【问题描述】:

我们有一个流程,用户可以通过该流程请求我们需要从源获取的文件。此来源不是最可靠的,因此我们使用 Amazon SQS 实施了一个队列。我们将下载 URL 放入队列中,然后使用我们用 Go 编写的小应用程序对其进行轮询。这个应用程序只是检索消息,下载文件,然后将其推送到我们存储它的 S3。一旦所有这些都完成了,它会回调一个服务,该服务会向用户发送电子邮件,让他们知道文件已准备好。

最初我写这个是为了创建 n 个通道,然后为每个通道附加 1 个 go-routine,并让 go-routine 处于无限循环中。这样我可以确保我一次只处理固定数量的下载。

我意识到这不是应该使用通道的方式,如果我现在理解正确的话,实际上应该有一个通道接收 n 个 go-routines渠道。每个 go-routine 都处于无限循环中,等待消息,当它接收到数据时,它会处理数据,做它应该做的一切,当它完成时,它会等待下一条消息。这使我可以确保我一次只处理 n 个文件。我认为这是正确的做法。我相信这是扇出,对吧?

不需要需要做的是将这些流程重新合并在一起。下载完成后,它会回调一个远程服务,以便处理该过程的其余部分。该应用无需执行任何其他操作。

好的,代码如下:

func main() {
    queue, err := ConnectToQueue() // This works fine...
    if err != nil {
        log.Fatalf("Could not connect to queue: %s\n", err)
    }

    msgChannel := make(chan sqs.Message, 10)

    for i := 0; i < MAX_CONCURRENT_ROUTINES; i++ {
        go processMessage(msgChannel, queue)
    }

    for {
        response, _ := queue.ReceiveMessage(MAX_SQS_MESSAGES)

        for _, m := range response.Messages {
            msgChannel <- m
        }
    }
}

func processMessage(ch <-chan sqs.Message, queue *sqs.Queue) {
    for {
        m := <-ch
        // Do something with message m

        // Delete message from queue when we're done
        queue.DeleteMessage(&m)
    }
}

我在这附近的任何地方吗?我有 n 个运行 go-routines(其中MAX_CONCURRENT_ROUTINES = n),在循环中我们将继续将消息传递到单个通道。这是正确的方法吗?我需要关闭任何东西还是可以让它无限期地运行?

我注意到的一件事是 SQS 正在返回消息,但是一旦我将 10 条消息传递到 processMessage()(10 是通道缓冲区的大小),实际上就没有进一步的消息被处理。

谢谢大家

【问题讨论】:

    标签: go concurrency channel goroutine


    【解决方案1】:

    看起来不错。几点说明:

    1. 您可以通过限制生成的工作例程数量以外的方式限制工作并行性。例如,您可以为收到的每条消息创建一个 goroutine,然后让生成的 goroutine 等待限制并行度的信号量。当然也有取舍,但您并不仅限于您所描述的方式。

      sem := make(chan struct{}, n)
      work := func(m sqs.Message) {
          sem <- struct{}{} // When there's room we can proceed
          // do the work
          <-sem // Free room in the channel
      }()
      for _, m := range queue.ReceiveMessage(MAX_SQS_MESSAGES) {
          for _, m0 := range m {
              go work(m0)
          }
      }
      
    2. 仅处理 10 条消息的限制是由堆栈中的其他地方引起的。可能您看到前 10 个填满通道的比赛,然后工作没有完成,或者您可能不小心从工作程序中返回。如果您的员工按照您描述的模式坚持不懈,那么您需要确定他们不会返回。

    3. 目前尚不清楚您是否希望进程在处理了一定数量的消息后返回。如果您确实希望退出此过程,则需要等待所有工作人员完成他们当前的任务,然后可能会发出信号让他们返回。查看sync.WaitGroup 以同步他们的完成,并有另一个频道来表示没有更多工作,或者关闭msgChannel,并在您的工作人员中处理。 (看看 2 元组返回通道接收表达式。)

    【讨论】:

    • 谢谢@matt-joiner,我当然已经回来了……我以前每条消息都有一个例行程序,完成后他们会返回。当我将其移至仅运行 10 时,我忘记将 return 更改为 continue
    • 好的,我已经处理完您的回复,现在我已经修复了return/continue 问题。感谢有关信号量的建议;我将阅读有关该主题的内容。我不需要这些工人回来,不。他们应该只处理包含回调 URI 的消息,所以他们只是调用它然后等待下一条消息。然而,感谢sync.WaitGroup 上的指针和 2 元组返回通道接收表达式。再次感谢您的帮助
    猜你喜欢
    • 2015-12-25
    • 2016-04-01
    • 1970-01-01
    • 1970-01-01
    • 2019-03-09
    • 2014-02-17
    • 2015-03-05
    • 1970-01-01
    相关资源
    最近更新 更多