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
相关产品推荐
相关产品推荐

