当前位置: > > > > 多个生产者,单个消费者:所有 goroutine 都在睡觉 – 死锁
多个生产者,单个消费者:所有 goroutine 都在睡觉 – 死锁
来源:stackoverflow
2024-04-28 21:51:35
0浏览
收藏
亲爱的编程学习爱好者,如果你点开了这篇文章,说明你对《多个生产者,单个消费者:所有 goroutine 都在睡觉 – 死锁》很感兴趣。本篇文章就来给大家详细解析一下,主要介绍一下,希望所有认真读完的童鞋们,都有实质性的提高。
问题内容
在继续工作之前,我一直遵循检查频道中是否有任何内容的模式:
func consume(msg <-chan message) {
for {
if m, ok := <-msg; ok {
fmt.println("more messages:", m)
} else {
break
}
}
}
这是基于这个视频的。这是我的完整代码:
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。如何使此代码工作以及为什么会出现此死锁错误(大概此模式应该有效?)。
解决方案
鉴于您在单个通道上确实有多个编写器,您会遇到一些挑战,因为在 go 中执行此操作的简单方法通常是在单个通道上有一个编写器,然后让该单个编写器作者在发送最后一个数据后关闭通道:
func produce(... args including channel) {
defer close(ch)
for stuff_to_produce {
ch <- item
}
}
此模式有一个很好的特性,无论您如何退出 produce,通道都会关闭,从而发出生产结束的信号。
您没有使用这种模式 – 您将一个通道传递给许多 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 完成的。整个生产者现在等待它们全部完成,然后在退出时关闭通道。
(如果将内部 productinparallel 函数编写为闭包,则无需显式向其传递 ch 和 pg。您也可以将 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 构造。)
另请注意,receive 的 wg(这可能是不必要的,具体取决于您如何构造其他事物)与生产者的等待组 pg 完全独立。虽然确实如所写的,在所有生产者完成之前消费者无法完成,但我们希望独立等待生产者完成,以便我们可以关闭整体生产者包装器中的通道。 p>
尝试这个代码,我做了一些修复使其工作:
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
// }
}
func receive(msg <-chan message, wg *sync.WaitGroup) {
for {
if m, ok := <-msg; ok {
fmt.Println("Received:", m)
wg.Done()
}
}
}
func consume(msg <-chan message) {
for {
if m, ok := <-msg; ok {
fmt.Println("More messages:", m)
} else {
break
}
}
}
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)
}
go receive(msgC, &wg)
wg.Wait()
fmt.Println("Breaking from receiving")
}
以上就是本文的全部内容了,是否有顺利帮助你解决问题?若是能给你带来学习上的帮助,请大家多多支持米云!更多关于Golang的相关知识,也可关注米云公众号。
