【问题标题】:Go worker pool while limiting number of goroutines and timeout for calculationsGo 工作池,同时限制 goroutine 的数量和计算超时
【发布时间】:2021-03-16 02:45:54
【问题描述】:

我有一个函数应该最多生成 N 个 goroutine,然后每个 goroutine 将从作业通道读取并进行一些计算。但是需要注意的是,如果计算花费的时间超过 X,请结束该计算并继续进行下一个计算。

func doStuff(){
    rules := []string{
        "a",
        "b",
        "c",
        "d",
        "e",
        "f",
        "g",
    }
    var (
        jobs    = make(chan []string, len(rules))
        res     = make(chan bool, len(rules))
        matches []string
    )

    w := func(jobs <-chan []string, results chan<- bool) {
        for j := range jobs {
            k, id := j[0], j[1]
            if id == "c" || id == "e" {
                time.Sleep(time.Second * 5)
            }
            m := match(k, id)
            res <- m
        }
    }
    N := 2
    for i := 0; i < N; i++ {
        go w(jobs, res)
    }

    for _, rl := range rules {
        jobs <- []string{"a", rl}
    }
    close(jobs)

    for i := 0; i < len(rules); i++ {
        select {
        case match := <-res:
            matches = append(matches, match)
        case <-time.After(time.Second):
        }
    }
    fmt.Println(matches)
}

预期结果是:

[a, b, d, f, g]

但我得到的是:

[a, b, d]

由于睡眠,似乎在其中一个 goroutine 可以完全完成之前从结果通道中读取结束。所以我添加了一个带有截止日期的上下文,但现在它无限期挂起:

    w := func(jobs <-chan []string, results chan<- string) {
        for j := range jobs {
            ctx, c := context.WithDeadline(context.Background(), time.Now().Add(time.Second*2))
            defer c()
            k, id := j[0], j[1]
            if id == "c" || id == "e" {
                time.Sleep(time.Second * 5)
            }
            m := match(k, id)
            select {
            case res <- m:
            case <-ctx.Done():
                fmt.Println("Canceled by timeout")
                continue
            }
        }
    }

我已经阅读了其他关于在发生超时时完全杀死 goroutine 的问题,但找不到任何关于超时时跳过的内容。

【问题讨论】:

    标签: go


    【解决方案1】:

    我为这样的用例制作了一个包。请查看此存储库:github.com/MicahParks/ctxerrgroup。

    这是一个完整示例,说明您的代码在使用包和流式传输结果时的外观。流式方法更节省内存。原始方法在最后打印之前将所有结果保存在内存中。

    package main
    
    import (
        "context"
        "log"
        "time"
    
        "github.com/MicahParks/ctxerrgroup"
    )
    
    func main() {
    
        // The number of worker goroutines to use.
        workers := uint(2)
    
        // Create an error handler that logs all errors.
        //
        // The original work item didn't return an error, so this is not required.
        var errorHandler ctxerrgroup.ErrorHandler
        errorHandler = func(_ ctxerrgroup.Group, err error) {
            log.Printf("A job in the worker pool failed.\nError: %s", err.Error())
        }
    
        // Create the group of workers.
        group := ctxerrgroup.New(workers, errorHandler)
    
        // Create the question specific assets.
        rules := []string{
            "a",
            "b",
            "c",
            "d",
            "e",
            "f",
            "g",
        }
        results := make(chan bool)
    
        // Create a parent timeout.
        timeout := time.Second
        parentTimeout, parentCancel := context.WithTimeout(context.Background(), timeout)
        defer parentCancel()
    
        // Iterate through all the rules to use.
        for _, rule := range rules {
    
            // Create a child context for this specific work item.
            ctx, cancel := context.WithCancel(parentTimeout)
    
            // Create and add the work item.
            group.AddWorkItem(ctx, cancel, func(workCtx context.Context) (err error) {
    
                // Deliberately shadow the rule so the next iteration doesn't take over.
                rule := rule
    
                // Do the work using the workCtx.
                results <- match(workCtx, "a", rule)
    
                return nil
            })
        }
    
        // Launch a goroutine that will close the results channel when everyone is finished.
        go func() {
            group.Wait()
            close(results)
        }()
    
        // Print the matches as the happen. This will not hang.
        for result := range results {
            log.Println(result)
        }
    
        // Wait for the group to finish.
        //
        // This is not required, but doesn't hurt as group.Wait is idempotent. It's here in case you remove the goroutine
        // waiting and closing the channel above.
        group.Wait()
    }
    
    // match is a function from the original question. It now accepts and properly uses the context argument.
    func match(ctx context.Context, key, id string) bool {
        panic("implement me")
    }
    

    【讨论】:

    • 嗨,所以我查看了 repo,我的理解是在创建有限的工作人员之后,如果工作人员完成工作或在通过 kill 通道发送死亡信号后,您将释放工作人员。所以本质上,我们只能做类似“它超时了,所以让我们释放工人,让它做下一件事情”之类的事情。我是否理解没有办法在执行过程中杀死一个 goroutine?
    • 您可以在执行期间通过使用其关联的cancel 类型的context.CancelFunc 函数来终止work item。这仅在您的 worker function 正确使用其上下文时才有效。 (work item和worker function的含义请参见README.md。)
    • Am I to understand there's no way to kill a goroutine in the middle of its execution? nop。您不能从外部处理程序执行此操作(请注意,go rutine 不会公开任何类型的 ID 或blog.sgmansfield.com/2015/12/goroutine-ids)。不过,您可以使用上下文或某些安全处理程序向它发出退出或返回的信号。这意味着在您的计算过程中,您必须检查上下文取消才能最终返回。这里没有魔法。
    • 啊,我明白了,所以我能做的最好的就是等待超时,如果它在阻塞通道的限制内,则启动另一个 goroutine。否则通道会被阻塞,我必须等到另一个 goroutine 完成并解除阻塞。
    【解决方案2】:

    所以 Micah 的答案肯定有效,但我想出了这个,它不需要第三方库。我刚刚修改了带有上下文的版本并提前移动了上下文检查:

    worker := func (ctx context.Context, wg *sync.WaitGroup, input <- chan string, res chan <- bool){
        defer wg.Done()
        for id := range input{
            if id == "c" || id == "e" {
                time.Sleep(time.Second * 5)
            }
            select{
            case: <- ctx.Done():
            return
            default:
            }
            // logic
            res <- match()
        }
    }
    
    rules := []string{"a", "b", "c", "d", "e"}
    input := make(chan string, len(rules))
    res := make(chan bool, len(rules))
    N := 2 // Limit # of goroutines
    
    var wg sync.WaitGroup
    ctx, cancel := context.WithTimeout(context.Background, time.Second * 3)
    defer cancel()
    for i := 0; i < N; i++{
        go worker(ctx, &wg, input, res)
    }
    
    for _, r := range rules{
        input <- r
    }
    close(r)
    
    wg.Wait()
    close(res)
    for r := range res{
        // do stuff
    }
    
    

    此答案基于 mh-cbon 的评论:

    不过,您可以使用上下文或某些安全处理程序向它发出退出或返回的信号。这意味着在您的计算过程中,您必须检查上下文取消才能最终返回。

    因此,我们不是使用通道来跟踪在任何给定时间产生的 goroutines 的数量,而是预先产生限制,以便我们可以检查超时。

    【讨论】:

      猜你喜欢
      • 2021-12-13
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2021-09-04
      • 2017-09-25
      • 2014-02-26
      • 2015-08-20
      • 1970-01-01
      相关资源
      最近更新 更多