Go语言生产者消费者模式死锁问题排查与修复求助
Go生产者消费者模式死锁问题修复方案
我是Go并发编程新手,写了一个基于goroutine和channel的生产者消费者示例:生产者持续生成随机字符串,消费者转成大写,期望运行2秒后停止。但程序只输出一个元素就触发死锁,报错fatal error: all goroutines are asleep - deadlock!,代码和运行输出如下:
package main import ( "fmt" "math/rand" "strings" "time" ) func producer(x []string, c chan string) { i := 1 for i > 0 { randomIndex := rand.Intn(len(x)) pick := x[randomIndex] c <- pick } } func consumer(x string, c chan string) { x1 := strings.ToUpper(x) c <- x1 } func main() { s := []string{"one", "two", "three", "four"} c1 := make(chan string) d1 := make(chan string) go producer(s, c1) go consumer(<-c1, d1) stop := time.After(2000 * time.Millisecond) for { select { case <-stop: fmt.Println("STOP AFTER 2 SEC!") return default: fmt.Println(<-d1) time.Sleep(50 * time.Millisecond) } } }
运行输出:
TWO fatal error: all goroutines are asleep - deadlock! goroutine 1 [chan receive]: main.main() goroutine 6 [chan send]: main.producer({0xc00004e040, 0x4, 0x0?}, 0x0?) created by main.main exit status 2
问题原因及修改方案
1. 消费者仅执行一次,启动逻辑错误
当前go consumer(<-c1, d1)是在主goroutine中先从c1取出一个值,再将该值传给consumer goroutine。consumer处理完这个值后就直接退出,无法处理后续生产者生成的消息。
修改方式:让consumer持续监听输入channel的消息,循环处理:
func consumer(cin chan string, cout chan string) { defer close(cout) // 处理完所有消息后关闭输出channel for x := range cin { x1 := strings.ToUpper(x) cout <- x1 } }
启动时直接传入channel,不再提前取值:go consumer(c1, d1)
2. 生产者无退出逻辑,导致阻塞
生产者的for i>0是无限循环,即使主程序收到停止信号,生产者仍会尝试往c1发送消息,最终因无人接收导致阻塞死锁。
修改方式:用context控制生产者生命周期,收到停止信号后关闭输入channel,通知consumer停止:
func producer(ctx context.Context, x []string, c chan string) { rand.Seed(time.Now().UnixNano()) // 初始化随机数种子,避免生成重复序列 for { select { case <-ctx.Done(): close(c) return default: randomIndex := rand.Intn(len(x)) pick := x[randomIndex] // 发送时也检查停止信号,避免阻塞 select { case c <- pick: case <-ctx.Done(): close(c) return } } } }
3. 主程序循环逻辑错误
主程序中select的default分支会强制从d1取值,当d1无消息时会直接阻塞,触发死锁。需将接收d1的操作放到select的case中,配合停止信号处理。
修改方式:调整主循环逻辑,用context通知生产者停止,并消费完剩余消息:
func main() { s := []string{"one", "two", "three", "four"} c1 := make(chan string) d1 := make(chan string) ctx, cancel := context.WithCancel(context.Background()) defer cancel() go producer(ctx, s, c1) go consumer(c1, d1) stop := time.After(2000 * time.Millisecond) for { select { case <-stop: fmt.Println("STOP AFTER 2 SEC!") cancel() // 通知生产者停止 // 消费剩余未处理的消息 for msg := range d1 { fmt.Println(msg) } return case msg, ok := <-d1: if !ok { // d1被关闭,说明consumer已退出 return } fmt.Println(msg) time.Sleep(50 * time.Millisecond) } } }
完整修正代码
package main import ( "context" "fmt" "math/rand" "strings" "time" ) func producer(ctx context.Context, x []string, c chan string) { rand.Seed(time.Now().UnixNano()) for { select { case <-ctx.Done(): close(c) return default: randomIndex := rand.Intn(len(x)) pick := x[randomIndex] select { case c <- pick: case <-ctx.Done(): close(c) return } } } } func consumer(cin chan string, cout chan string) { defer close(cout) for x := range cin { x1 := strings.ToUpper(x) cout <- x1 } } func main() { s := []string{"one", "two", "three", "four"} c1 := make(chan string) d1 := make(chan string) ctx, cancel := context.WithCancel(context.Background()) defer cancel() go producer(ctx, s, c1) go consumer(c1, d1) stop := time.After(2000 * time.Millisecond) for { select { case <-stop: fmt.Println("STOP AFTER 2 SEC!") cancel() for msg := range d1 { fmt.Println(msg) } return case msg, ok := <-d1: if !ok { return } fmt.Println(msg) time.Sleep(50 * time.Millisecond) } } }
内容的提问来源于stack exchange,提问作者DvdiidI
相关产品推荐
相关产品推荐

