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

Go生产者消费者模型死锁风险分析及规避方案

Go生产者消费者模型的死锁风险与问题修复

问题背景

我实现了一个多生产者多消费者的Go并发模型,核心特性包括:

  • 多生产者、多消费者共享同一数据通道
  • 错误处理机制:任一生产者或消费者出错时,所有工作协程都会被终止

我原本担心会出现“消费者全部关闭但生产者仍向通道写入”的死锁场景,因此在写入数据前添加了Context检查,但仍有疑问:会不会出现“生产者检查Context未取消后,Context突然被取消导致消费者关闭,随后生产者尝试写入通道引发死锁”的情况?如果存在该风险,该如何规避?

核心代码

package main

import (
    "context"
    "fmt"
    "sync"
)

func main() {
    a1 := []int{1, 2, 3, 4, 5}
    a2 := []int{5, 4, 3, 1, 1}
    a3 := []int{6, 7, 8, 9}
    a4 := []int{1, 2, 3, 4, 5}
    a5 := []int{5, 4, 3, 1, 1}
    a6 := []int{6, 7, 18, 9}
    arrayOfArray := [][]int{a1, a2, a3, a4, a5, a6}

    ctx, cancel := context.WithCancel(context.Background())
    ch1 := read(ctx, arrayOfArray)

    messageCh := make(chan int)
    errCh := make(chan error)

    producerWg := &sync.WaitGroup{}
    for i := 0; i < 3; i++ {
        producerWg.Add(1)
        producer(ctx, producerWg, ch1, messageCh, errCh)
    }

    consumerWg := &sync.WaitGroup{}
    for i := 0; i < 3; i++ {
        consumerWg.Add(1)
        consumer(ctx, consumerWg, messageCh, errCh)
    }

    firstError := handleAllErrors(ctx, cancel, errCh)

    producerWg.Wait()
    close(messageCh)

    consumerWg.Wait()
    close(errCh)

    fmt.Println(<-firstError)
}

func read(ctx context.Context, arrayOfArray [][]int) <-chan []int {
    ch := make(chan []int)

    go func() {
        defer close(ch)

        for i := 0; i < len(arrayOfArray); i++ {
            select {
            case <-ctx.Done():
                return
            case ch <- arrayOfArray[i]:
            }
        }
    }()

    return ch
}

func producer(ctx context.Context, wg *sync.WaitGroup, in <-chan []int, messageCh chan<- int, errCh chan<- error) {
    go func() {
        defer wg.Done()
        for {
            select {
            case <-ctx.Done():
                return
            case arr, ok := <-in:
                if !ok {
                    return
                }

                for i := 0; i < len(arr); i++ {
                    // simulating an error.
                    //if arr[i] == 10 {
                    //  errCh <- fmt.Errorf("producer interrupted")
                    //}

                    select {
                    case <-ctx.Done():
                        return
                    case messageCh <- 2 * arr[i]:
                    }
                }
            }
        }
    }()
}

func consumer(ctx context.Context, wg *sync.WaitGroup, messageCh <-chan int, errCh chan<- error) {
    go func() {
        wg.Done()

        for {
            select {
            case <-ctx.Done():
                return
            case n, ok := <-messageCh:
                if !ok {
                    return
                }
                fmt.Println("consumed: ", n)

                // simulating erros
                //if n == 10 {
                //  errCh <- fmt.Errorf("output error during write")
                //}
            }
        }
    }()
}

func handleAllErrors(ctx context.Context, cancel context.CancelFunc, errCh chan error) <-chan error {
    firstErrCh := make(chan error, 1)
    isFirstError := true
    go func() {
        defer close(firstErrCh)
        for err := range errCh {
            select {
            case <-ctx.Done():
            default:
                cancel()
            }
            if isFirstError {
                firstErrCh <- err
                isFirstError = !isFirstError
            }
        }
    }()

    return firstErrCh
}

问题解答

1. 你担忧的死锁场景不会在当前代码中发生

你在写入messageCh前使用的是如下写法:

select {
case <-ctx.Done():
    return
case messageCh <- 2 * arr[i]:
}

Go的select语句是原子性执行的:它会同时监听所有case,哪个case先就绪就执行哪个分支,不存在“检查完ctx未取消,之后ctx才被取消”的时间窗口。如果在写入前ctx被取消,<-ctx.Done()会先触发,协程直接返回,不会执行写入操作;如果写入操作先就绪,就会成功发送数据。因此不会出现你担心的死锁场景。

但如果是如下错误写法(分开检查和写入),才会存在竞态条件:

// 错误示例:存在竞态风险
if ctx.Err() == nil {
    messageCh <- 2 * arr[i] // 此处可能因ctx取消、无接收方而阻塞死锁
}

这种情况下,检查ctx和写入通道是两个独立操作,中间可能被ctx取消信号打断,导致写入时无接收方,最终阻塞死锁。

2. 当前代码存在的严重bug:WaitGroup使用错误

你的consumer函数中,协程启动后立即调用wg.Done(),这会导致main函数中的consumerWg.Wait()立即完成,随后执行close(messageCh)。此时如果生产者还在向messageCh写入数据,会直接引发panic(向已关闭的通道发送数据)。

修复方法:将wg.Done()改为延迟调用,确保协程完成所有任务后再通知WaitGroup:

func consumer(ctx context.Context, wg *sync.WaitGroup, messageCh <-chan int, errCh chan<- error) {
    go func() {
        defer wg.Done() // 改为defer,协程退出时才执行
        for {
            select {
            case <-ctx.Done():
                return
            case n, ok := <-messageCh:
                if !ok {
                    return
                }
                fmt.Println("consumed: ", n)

                // simulating erros
                //if n == 10 {
                //  errCh <- fmt.Errorf("output error during write")
                //}
            }
        }
    }()
}

3. 其他注意事项

  • 确保生产者全部退出后再关闭messageCh:当前代码中producerWg.Wait()后调用close(messageCh)的逻辑是正确的,能避免生产者向已关闭通道写入。
  • 错误通道errCh的关闭:main函数中在consumerWg.Wait()后关闭errCh,能保证所有错误发送操作都在通道关闭前完成,避免handleAllErrors中的range循环提前退出。

内容的提问来源于stack exchange,提问作者JavaDeveloper

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.03 22:10:40