【问题标题】:Distribute the same keyword to multiple goroutines将相同的关键字分配给多个 goroutine
【发布时间】: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


【解决方案1】:

假设我正确理解了这个问题:

恐怕你无能为力。您的资源有限(假设所有资源都是有限的),因此如果您输入的数据写入速度比您处理它的速度更快,则需要进行一些同步。最后,整个过程将与最慢的工人一样快。

如果您确实需要尽快从可用的工作人员那里获取数据,那么您能做的最好的事情就是添加某种缓冲。但是缓冲区的大小必须受到限制(即使您在云中运行,它也会受到您的钱包的限制),因此假设输入的洪流永无止境,它只会将阻塞推迟到将来的某个时间,您将开始看到“同步” ” 再次。

您在问题中提出的所有想法都是基于缓冲数据。即使您为每个关键字-工作者对运行一个例程,这也会在每个例程中缓冲一个元素,除非您实施例程总数的限制,否则您将耗尽内存。即使您总是为最快的工作人员留出一些空间来生成新的例程,输入队列也将无法交付新项目,因为它会被最慢的工作人员阻塞。

缓冲解决您的问题,如果平均而言您的输入比处理时间慢,但您偶尔会出现峰值。如果您的缓冲区足够大,您就可以适应吞吐量的增加,并且您最快的工作人员可能不会注意到任何事情。

解决方案?

由于 go 带有缓冲通道,这是最容易实现的(icza 在评论中也建议)。只需给每个工人一个缓冲区。如果你知道哪个工人最慢,你可以给它一个更大的缓冲区。在这种情况下,您会受到机器内存的限制。

如果您对单机内存限制不满意,那么可以,根据您的想法,您可以“简单”地将每个工作人员的缓冲区(队列)存储在硬盘上。但这也是有限的,只是将阻塞情况推迟到以后。这与您的 Amazon SQS 提案基本相同(您可以在云中保留缓冲区,但您需要合理限制它或为账单做准备。)

最后一点,根据您正在构建的系统,缓冲如此大规模的项目可能不是一个好主意,从而为较慢的工作人员建立积压 - 通常不希望有工作人员输入流之后的几小时、几天、几周,这就是无限缓冲区会发生的情况。那么真正的答案是:提高你最慢的工人更快地处理事情。 (并添加 一些 缓冲以改善体验。)

【讨论】:

  • 有点希望我错过了一些愚蠢的明显解决方案,就像我在这里提问时通常会遇到的情况一样。加速慢速工作人员的问题在于它受到 http 请求的限制,我想我可以扔一些更多的 goroutine 来尝试给予它更多的优先级,看看这有什么帮助。
猜你喜欢
  • 2014-07-06
  • 1970-01-01
  • 2015-02-09
  • 1970-01-01
  • 1970-01-01
  • 2023-03-10
  • 1970-01-01
  • 2019-01-25
  • 1970-01-01
相关资源
最近更新 更多