【发布时间】:2015-08-26 22:05:31
【问题描述】:
我有类似这个模拟(下面的代码)的东西,它将相同的关键字分发给多个 goroutine,除了 goroutine 都花费不同的时间来处理关键字但可以彼此独立运行,因此它们不需要任何同步。下面给出的分发解决方案清楚地同步了 goroutine。
我只是想把这个想法抛诸脑后,看看其他人会如何处理这种类型的分发,因为我认为它相当普遍,并且其他人之前已经考虑过。
以下是我想出的其他一些解决方案,以及为什么它们对我来说似乎有点不对劲:
每个关键字一个 goroutine
每次出现新关键字时都会产生一个 goroutine 来处理分布
为要更新的每个 goroutine 给关键字一个位掩码或其他东西
这样,一旦所有的工人都触及了关键字,它就可以被删除,我们可以继续前进
给每个工人自己的堆栈来工作
这似乎有点吸引人,只需给每个工作人员一个堆栈来添加每个关键字,但我们最终会遇到一个问题,即计划运行这么长时间,会占用大量内存
所有这些的问题是我的代码应该运行很长时间,无人看管,这将导致关键字或 goroutine 的大量堆积,因为懒惰的工作人员比其他工作人员花费的时间更长。为每个工作人员提供自己的 Amazon SQS 队列或自己实现类似的东西似乎会很好。
编辑:
将关键字存储在程序外
我只是想这样做,我也许可以将关键字存储在程序之外,直到他们都抓住它,然后在它用完后将其删除。实际上,这对我来说没问题,我没有用完磁盘空间的问题
无论如何,这里是等待所有人完成的方法的示例:
package main
import (
"flag"
"fmt"
"math/rand"
"os"
"os/signal"
"strconv"
"time"
)
var (
shutdown chan struct{}
count = flag.Int("count", 5, "number to run")
)
type sleepingWorker struct {
name string
sleep time.Duration
ch chan int
}
func NewQuicky(n string) sleepingWorker {
var rq sleepingWorker
rq.name = n
rq.ch = make(chan int)
rq.sleep = time.Duration(rand.Intn(5)) * time.Second
return rq
}
func (r sleepingWorker) Work() {
for {
fmt.Println(r.name, "is about to sleep, number:", <-r.ch)
time.Sleep(r.sleep)
}
}
func NewLazy() sleepingWorker {
var rq sleepingWorker
rq.name = "Lazy slow worker"
rq.ch = make(chan int)
rq.sleep = 20 * time.Second
return rq
}
func distribute(gen chan int, workers ...sleepingWorker) {
for kw := range gen {
for _, w := range workers {
fmt.Println("sending keyword to:", w.name)
select {
case <-shutdown:
return
case w.ch <- kw:
fmt.Println("keyword sent to:", w.name)
}
}
}
}
func main() {
flag.Parse()
shutdown = make(chan struct{})
go func() {
c := make(chan os.Signal, 1)
signal.Notify(c, os.Interrupt)
<-c
close(shutdown)
}()
x := make([]sleepingWorker, *count)
for i := 0; i < (*count)-1; i++ {
x[i] = NewQuicky(strconv.Itoa(i))
go x[i].Work()
}
x[(*count)-1] = NewLazy()
go x[(*count)-1].Work()
gen := make(chan int)
go distribute(gen, x...)
go func() {
i := 0
for {
i++
select {
case <-shutdown:
return
case gen <- i:
}
}
}()
<-shutdown
os.Exit(0)
}
【问题讨论】:
-
为工人使用缓冲通道有什么问题?这样它们就不会“同步”(直到缓冲区被填满,这将是“吞吐量”)。例如。
rq.ch = make(chan int, 20) -
缓冲通道只会将同步向后移动一点,并不会真正改变任何东西。但是,是的,这是可能的。
-
offtopic:在您的实际实施中,您可能希望使用
sync.WaitGroup来确保您的每个工人都已完成,以防他们需要清理。我相信示例代码,不能保证您的每个工作人员都会进入从shutdown接收的案例,因为一旦您通过主函数中的封闭通道,程序就会退出。 -
@ogc-nick 是的,我知道,这只是对功能的快速模拟,而不是实际的退出条件。实际代码使用等待组,不过谢谢!
标签: go