【问题标题】:Shutdown "worker" go routine after buffer is empty缓冲区为空后关闭“工人”执行例行程序
【发布时间】:2015-11-29 17:27:09
【问题描述】:

我希望我的日常工作人员(下面代码中的ProcessToDo())等到所有“排队”的工作都处理完后再关闭。

工作例程有一个“待办事项”通道(缓冲),工作通过该通道发送给它。它有一个“完成”通道来告诉它开始关闭。文档说,如果满足多个选择,则通道上的选择将选择一个“伪随机值”......这意味着在所有缓冲工作完成之前触发关闭(返回)。

在下面的代码示例中,我希望打印所有 20 条消息...

package main

import (
    "time"
    "fmt"
)


func ProcessToDo(done chan struct{}, todo chan string) {
    for {
        select {
        case work, ok := <-todo:
            if !ok {
                fmt.Printf("Shutting down ProcessToDo - todo channel closed!\n")
                return
            }
            fmt.Printf("todo: %q\n", work)
            time.Sleep(100 * time.Millisecond)
        case _, ok := <-done:
            if ok {
                fmt.Printf("Shutting down ProcessToDo - done message received!\n")
            } else {
                fmt.Printf("Shutting down ProcessToDo - done channel closed!\n")
            }
            close(todo)
            return
        }
    }
}

func main() {

    done := make(chan struct{})
    todo := make(chan string, 100)

    go ProcessToDo(done, todo)

    for i := 0; i < 20; i++ {
        todo <- fmt.Sprintf("Message %02d", i)
    }

    fmt.Println("*** all messages queued ***")
    time.Sleep(1 * time.Second)
    close(done)
    time.Sleep(4 * time.Second)
}

【问题讨论】:

    标签: go concurrency shutdown channel goroutine


    【解决方案1】:

    让通道的消费者关闭它通常是个坏主意,因为在关闭的通道上发送是一种恐慌。

    在这种情况下,如果您不想在发送所有消息之前中断消费者,只需使用for...range 循环并在完成后关闭通道。您还需要一个像 WaitGroup 这样的信号来等待 goroutine 完成(而不是使用 time.Sleep)

    http://play.golang.org/p/r97vRPsxEb

    var wg sync.WaitGroup
    
    func ProcessToDo(todo chan string) {
        defer wg.Done()
        for work := range todo {
            fmt.Printf("todo: %q\n", work)
            time.Sleep(100 * time.Millisecond)
    
        }
        fmt.Printf("Shutting down ProcessToDo - todo channel closed!\n")
    
    }
    
    func main() {
        todo := make(chan string, 100)
        wg.Add(1)
        go ProcessToDo(todo)
    
        for i := 0; i < 20; i++ {
            todo <- fmt.Sprintf("Message %02d", i)
        }
    
        fmt.Println("*** all messages queued ***")
        close(todo)
        wg.Wait()
    }
    

    【讨论】:

    • 感谢您回答我的问题。 Icza 是第一个进来的,所以我不得不接受他的回答。当然也支持你的。
    【解决方案2】:

    done 通道在您的情况下是完全没有必要的,因为您可以通过关闭 todo 通道本身来发出关闭信号。

    并在通道上使用for range,它将迭代直到通道关闭并且其缓冲区为空。

    您应该有一个done 频道,但前提是 goroutine 本身可以发出信号表明它已完成工作,因此主 goroutine 可以继续或退出。

    这个变体等同于你的变体,更简单,不需要time.Sleep() 调用来等待其他goroutines(无论如何这都会太错误和不确定)。试试Go Playground

    func ProcessToDo(done chan struct{}, todo chan string) {
        for work := range todo {
            fmt.Printf("todo: %q\n", work)
            time.Sleep(100 * time.Millisecond)
        }
        fmt.Printf("Shutting down ProcessToDo - todo channel closed!\n")
        done <- struct{}{} // Signal that we processed all jobs
    }
    
    func main() {
        done := make(chan struct{})
        todo := make(chan string, 100)
    
        go ProcessToDo(done, todo)
    
        for i := 0; i < 20; i++ {
            todo <- fmt.Sprintf("Message %02d", i)
        }
    
        fmt.Println("*** all messages queued ***")
        close(todo)
        <-done // Wait until the other goroutine finishes all jobs
    }
    

    还要注意,worker goroutines 应该使用defer 发出完成信号,这样主goroutine 就不会在等待worker 以某种意外方式返回或出现恐慌时卡住。所以它应该像这样开始:

    defer func() {
        done <- struct{}{} // Signal that we processed all jobs
    }()
    

    您也可以使用sync.WaitGroup 将主 goroutine 同步到 worker(等待它)。事实上,如果你计划使用多个工作 goroutine,这比从 done 通道读取多个值更干净。此外,使用WaitGroup 发出完成信号更简单,因为它有一个Done() 方法(这是一个函数调用),因此您不需要匿名函数:

    defer wg.Done()
    

    有关WaitGroup 的完整示例,请参阅JimB's anwser

    如果您想使用多个工作 goroutine,使用 for range 也是惯用的:通道是同步的,因此您不需要任何额外的代码来同步对 todo 通道的访问或从它接收的作业。如果你关闭main() 中的todo 通道,这将正确地向所有工作goroutines 发出信号。但当然,所有排队的作业都将被接收和处理一次。

    现在采用使用WaitGroup 的变体使主goroutine 等待worker(JimB 的回答):如果你想要多个worker goroutine 怎么办;并发处理您的工作(并且很可能是并行处理)?

    您需要在代码中添加/更改的唯一内容是:真正启动其中的多个:

    for i := 0; i < 10; i++ {
        wg.Add(1)
        go ProcessToDo(todo)
    }
    

    无需更改任何其他内容,您现在拥有一个正确的并发应用程序,该应用程序使用 10 个并发 goroutines 接收和处理您的作业。而且我们没有使用任何“丑陋的”time.Sleep()(我们使用了一个但只是为了模拟慢速处理,而不是等待其他 goroutine),并且您不需要任何额外的同步。

    【讨论】:

    • 太棒了!谢谢。我使用 Go 的次数越多,了解的越多,我就越喜欢它。非常优雅,非常简单。再次感谢您非常彻底和有帮助的回答。
    【解决方案3】:

    我认为接受的答案对于这个特定示例非常有效。然而,要回答“缓冲区为空后关闭“工人”执行例行程序”的问题 - 可能有更优雅的解决方案。

    worker 可以在缓冲区为空时返回,而无需通过关闭通道来发出信号。

    如果工作人员需要处理的任务数量未知,这将特别有用。

    在这里查看:https://play.golang.org/p/LZ1y0eIRMeS

    package main
    
    import (
        "fmt"
        "time"
        "math/rand"
    )
    
    func main() {
        rand.Seed(time.Now().UnixNano())
        ch := make(chan interface{}, 10)
    
        go worker(ch)
        for i := 1; i <= rand.Intn(9) + 1; i++ {
                ch <- i
        }
    
        blocker := make(chan interface{})
        <-blocker
    }
    
    func worker(ch chan interface{}){   
        for {
            select {
            case msg := <- ch:
                fmt.Println("msg: ", msg)
            default:
                fmt.Println("exiting worker")
                return
            }
        }       
    }
    

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2023-03-11
      • 1970-01-01
      • 2019-07-29
      • 1970-01-01
      • 2022-01-02
      • 2016-02-19
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多