gpt4 book ai didi

go - 多个生产者,单个消费者 : all goroutines are asleep - deadlock

转载 作者:行者123 更新时间:2023-12-03 02:23:22 26 4
gpt4 key购买 nike

在继续工作之前,我一直遵循检查 channel 中是否有任何内容的模式:

func consume(msg <-chan message) {
for {
if m, ok := <-msg; ok {
fmt.Println("More messages:", m)
} else {
break
}
}
}

基于此video 。这是我的完整代码:

package main

import (
"fmt"
"strconv"
"strings"
"sync"
)

type message struct {
body string
code int
}

var markets []string = []string{"BTC", "ETH", "LTC"}

// produces messages into the chan
func produce(n int, market string, msg chan<- message, wg *sync.WaitGroup) {
// for i := 0; i < n; i++ {
var msgToSend = message{
body: strings.Join([]string{"market: ", market, ", #", strconv.Itoa(1)}, ""),
code: 1,
}
fmt.Println("Producing:", msgToSend)
msg <- msgToSend
// }
wg.Done()
}

func receive(msg <-chan message, wg *sync.WaitGroup) {
for {
if m, ok := <-msg; ok {
fmt.Println("Received:", m)
} else {
fmt.Println("Breaking from receiving")
break
}
}
wg.Done()
}

func main() {
wg := sync.WaitGroup{}
msgC := make(chan message, 100)
defer func() {
close(msgC)
}()
for ix, market := range markets {
wg.Add(1)
go produce(ix+1, market, msgC, &wg)
}
wg.Add(1)
go receive(msgC, &wg)
wg.Wait()
}

如果你尝试运行它,我们会在打印即将打破的消息之前陷入僵局。说实话,这是有道理的,因为上次,当 chan 中没有其他内容时,我们试图提取该值,所以我们得到了这个错误。但这种模式是行不通的 if m, ok := <- msg; ok 。如何使此代码工作以及为什么会出现此死锁错误(大概此模式应该有效?)。

最佳答案

鉴于您在单个 channel 上确实有多个编写器,您会遇到一些挑战,因为在 Go 中执行此操作的简单方法通常是在单个 channel 上有一个编写器,然后让该单个编写器作者在发送最后一个数据后关闭 channel :

func produce(... args including channel) {
defer close(ch)
for stuff_to_produce {
ch <- item
}
}

这个模式有一个很好的特性,无论你如何摆脱 produce , channel 关闭,表示生产结束。

您没有使用这种模式 - 您将一个 channel 传递给许多 goroutine,每个 goroutine 都可以发送一条消息 - 因此您需要移动 close (或者,当然,使用其他一些模式)。表达您需要的模式的最简单方法是:

func overall_produce(... args including channel ...) {
var pg sync.WaitGroup
defer close(ch)
for stuff_to_produce {
pg.Add(1)
go produceInParallel(ch, &pg) // add more args if appropriate
}
pg.Wait()
}

pg计数器累计活跃生产者。每人必须调用pg.Done()表明它是使用 ch 完成的。整个生产者现在等待它们全部完成,然后在退出时关闭 channel 。

(如果将内部 produceInParallel 函数编写为闭包,则无需显式向其传递 chpg。您也可以将 overallProducer 编写为闭包。)

请注意,您的单个消费者的循环可能最好使用 for ... range 来表达。构造:

func receive(msg <-chan message, wg *sync.WaitGroup) {
for m := range msg {
fmt.Println("Received:", m)
}
wg.Done()
}

(您提到了将 select 添加到循环中的意图,以便在消息尚未准备好时可以执行其他计算。如果该代码无法分离到独立的 goroutine 中,那么您实际上需要爱好者 m, ok := <-msg 构造。)

另请注意 wg对于 receive - 这可能是不必要的,具体取决于你如何构建其他事物 - 完全独立于 WaitGroup pg对于生产者来说。虽然确实如所写的,在所有生产者完成之前消费者无法完成,但我们希望独立等待生产者完成,以便我们可以关闭整体生产者包装器中的 channel 。

关于go - 多个生产者,单个消费者 : all goroutines are asleep - deadlock,我们在Stack Overflow上找到一个类似的问题: https://stackoverflow.com/questions/58793428/

26 4 0
Copyright 2021 - 2024 cfsdn All Rights Reserved 蜀ICP备2022000587号
广告合作:1813099741@qq.com 6ren.com