You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.11 01:25:27