实现Pipeline并发模式时遭遇死锁问题求助
问题原因分析
死锁的核心原因是第一个job的输出通道outs[1]没有被任何后续job读取,且该通道是无缓冲的:
- 第一个job会向
outs[1]写入7个元素,但第二个job(pipe函数)完全忽略了输入通道,根本不会从outs[1]读取数据。 - 无缓冲通道的发送操作必须等待接收者才能完成,因此第一个job在执行
out <- fibNum时会永久阻塞,对应的goroutine无法走到wg.Done(),最终导致wg.Wait()无限等待,触发死锁。
解决方案(无需修改main函数)
由于不能修改main函数,我们可以通过修改ExecutePipeline2来解决这个问题,有两种可行方案:
方案1:使用带足够缓冲的通道
将所有通道改为带足够缓冲的类型,确保发送操作不会因暂时无接收者而阻塞。缓冲大小只要能容纳job可能输出的最大数据量即可:
func ExecutePipeline2(jobs ...intJob) { outs := make([]chan int, len(jobs)+1) wg := sync.WaitGroup{} for i := 0; i < len(outs); i++ { // 创建带足够缓冲的通道,避免发送阻塞 outs[i] = make(chan int, 100) } for i, job := range jobs { job := job in, out := outs[i], outs[i+1] i := i wg.Add(1) go func() { job(in, out) fmt.Printf("job %d closed\n", i) close(out) wg.Done() }() } wg.Wait() }
方案2:为未被消费的通道启动消费goroutine
为每个输入通道启动一个后台goroutine,自动消费通道内的所有数据,确保发送方不会因无接收者阻塞:
func ExecutePipeline2(jobs ...intJob) { outs := make([]chan int, len(jobs)+1) wg := sync.WaitGroup{} for i := 0; i < len(outs); i++ { outs[i] = make(chan int) } // 为每个输入通道启动消费goroutine,避免发送方阻塞 for i := 0; i < len(jobs); i++ { ch := outs[i] go func() { // 消费通道内所有数据直到通道关闭 for range ch { } }() } for i, job := range jobs { job := job in, out := outs[i], outs[i+1] i := i wg.Add(1) go func() { job(in, out) fmt.Printf("job %d closed\n", i) close(out) wg.Done() }() } wg.Wait() }
验证效果
修改后运行代码,三个job都会正常执行完成:
- 第一个job写入7个元素到
outs[1],要么被缓冲容纳,要么被消费goroutine读取,不会阻塞; - 第二个job写入5个元素到
outs[2]; - 第三个job读取
outs[2]的5个元素并打印,最终所有goroutine完成,WaitGroup正常结束,无死锁。
内容的提问来源于stack exchange,提问作者SYKO
相关产品推荐
相关产品推荐

