【问题标题】:How to read from a single channel shared by multiple goroutines如何从多个 goroutine 共享的单个通道中读取
【发布时间】:2018-10-15 07:54:53
【问题描述】:

我有几个 goroutine 写入同一个频道。如果我使用缓冲通道,我可以检索输入。但是,如果使用无缓冲通道,我只能读取大约一半的值:

func testAsyncFunc2() {

    ch := make(chan int,10)

    fmt.Println("testAsyncFunc2")
    wg.Add(10)
    for  i :=0; i < 10; i++  {
        go sender3(ch, i)
        wg.Done()
    }

    receiver3(ch)

    close(ch)
    wg.Wait()
}

这是接收函数:

func receiver3(ch chan int) {
    for {
        select {
        case <-ch:
            fmt.Println(<-ch)
        default:
            fmt.Println("Done...")
            return
        }
    }
}

发件人功能:

func sender3(ch chan int, i int) {
    ch <- i
}

然后输出:

testAsyncFunc 2 0 4 6 8 2 完成...

虽然我希望得到 10 个数字。

【问题讨论】:

  • 显示 sender3 的代码。 select 语句中的接收分支接收两个值,但只打印一个。

标签: go


【解决方案1】:

选择默认值并不像您想象的那样工作。如果一个选择块有一个默认案例,如果其他案例都没有准备好读取,它将立即选择它。

这意味着你的receiver3很可能会达到默认情况。

【讨论】:

    【解决方案2】:

    如果您不创建缓冲区通道,代码将返回错误,原因是通道在发送所有值之前关闭。

    不要关闭通道并等待 go 例程完成。如果要关闭通道,请在接收器 go 例程中接收到通道上发送的所有值时关闭。

    fmt.Println("testAsyncFunc2")
    for  i :=0; i < 10; i++  {
        wg.Add(1)
        go sender3(ch, i)
    }
    
    receiver3(ch)
    close(ch) // this will close the channels before all the values sent on it will be received.
    wg.Wait()
    

    还有一点需要注意的是,当您在 for 循环中启动 go 例程后减小计数器时,您已将等待组计数器增加到 10,这是错误的。当发送者 go 例程完成执行时,您应该减少其内部的等待组计数器。

    func sender(ch chan int, int){
         defer wg.Done()
    }
    

    在 for 循环中,选择默认条件将在接收到通道上发送的所有值之前运行,这就是所有发送的值都不会打印的原因。因为当通道上没有发送值时循环将返回。

    func receiver3(ch chan int) {
        for {
            select {
            case <-ch:
                fmt.Println(<-ch)
            default: // this condition will run when value is not available on the channel.
                fmt.Println("Done...")
                return
            }
        }
    }
    

    创建一个 go 例程来关闭通道并等待发送者 go 例程完成。所以下面的代码将等待你所有的 goroutine 完成在通道上发送值,然后它会关闭通道:

    package main
    
    import (
        "fmt"
        "sync"
    )
    
    var wg sync.WaitGroup
    
    func main() {
        ch := make(chan int)
        fmt.Println("testAsyncFunc2")
        for i := 0; i < 10; i++ {
            wg.Add(1)
            go sender(ch, i)
        }
        receiver3(ch)
        go func() {
            defer close(ch)
            wg.Wait()
        }()
    }
    
    func receiver3(ch <-chan int) {
        for i := 0; i < 10; i++ {
            select {
            case value, ok := <-ch:
                if !ok {
                    ch = nil
                }
                fmt.Println(value)
            }
            if ch == nil {
                break
            }
        }
    }
    
    func sender(ch chan int, i int) {
        defer wg.Done()
        ch <- i
    }
    

    输出

    testAsyncFunc2
    9
    0
    1
    2
    3
    4
    5
    6
    7
    8
    

    Go playground 上的工作代码

    【讨论】:

      猜你喜欢
      • 2021-09-27
      • 1970-01-01
      • 2019-10-08
      • 1970-01-01
      • 2014-09-30
      • 2020-08-02
      • 1970-01-01
      • 2019-11-18
      • 2022-01-02
      相关资源
      最近更新 更多