【发布时间】: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