如何从多个 goroutine 共享的单个通道中读取数据

s95*_*163 5 go

我有几个 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()
}
Run Code Online (Sandbox Code Playgroud)

这是接收器函数:

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

发送功能:

func sender3(ch chan int, i int) {
    ch <- i
}
Run Code Online (Sandbox Code Playgroud)

和输出:

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

虽然我希望能拿回 10 个号码。

Him*_*shu 1

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

不要关闭通道并等待 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()
Run Code Online (Sandbox Code Playgroud)

还要注意的另一件事是,在 for 循环内启动 go 例程后减少计数器时,您已将等待组计数器增加到 10,这是错误的。当发送方 go 例程完成执行时,您应该减少它内的等待组计数器。

func sender(ch chan int, int){
     defer wg.Done()
}
Run Code Online (Sandbox Code Playgroud)

在 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
        }
    }
}
Run Code Online (Sandbox Code Playgroud)

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

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
}
Run Code Online (Sandbox Code Playgroud)

输出

testAsyncFunc2
9
0
1
2
3
4
5
6
7
8
Run Code Online (Sandbox Code Playgroud)

Go 游乐场上的工作代码