使用WaitGroups与Buffered Channels的Go代码死锁原因排查
WaitGroups、Buffered Channels与死锁
这段Go代码出现死锁,需求是通过inputChan发送数据,再从outputChan读取数据。
原始代码:
package main import ( "fmt" "sync" ) func listStuff(wg *sync.WaitGroup, workerID int, inputChan chan int, outputChan chan int) { defer wg.Done() for i := range inputChan { fmt.Println("sending ", i) outputChan <- i } } func List(workers int) ([]int, error) { _output := make([]int, 0) inputChan := make(chan int, 1000) outputChan := make(chan int, 1000) var wg sync.WaitGroup wg.Add(workers) fmt.Printf("+++ Spinning up %v workers\n", workers) for i := 0; i < workers; i++ { go listStuff(&wg, i, inputChan, outputChan) } for i := 0; i < 3000; i++ { inputChan <- i } done := make(chan struct{}) go func() { close(done) close(inputChan) close(outputChan) wg.Wait() }() for o := range outputChan { fmt.Println("reading from channel...") _output = append(_output, o) } <-done fmt.Printf("+++ output len: %v\n", len(_output)) return _output, nil } func main() { List(5) }
死锁原因
- 双向阻塞:主goroutine向
inputChan发送3000个数据,当inputChan缓冲(1000)填满后,需要等待worker取走数据才能继续发送;而worker从inputChan取数据后向outputChan发送,当outputChan缓冲(1000)填满后,worker阻塞,无法继续读取inputChan,导致主goroutine也阻塞,形成死锁。 - 通道关闭顺序错误:匿名goroutine提前关闭
outputChan,不仅会导致worker向已关闭通道发送数据触发panic,还无法保证worker完成所有数据转发。
修复代码
package main import ( "fmt" "sync" ) func listStuff(wg *sync.WaitGroup, workerID int, inputChan chan int, outputChan chan int) { defer wg.Done() for i := range inputChan { fmt.Println("worker", workerID, "sending", i) outputChan <- i } } func List(workers int) ([]int, error) { _output := make([]int, 0) inputChan := make(chan int, 1000) outputChan := make(chan int, 1000) var wg sync.WaitGroup wg.Add(workers) fmt.Printf("+++ Spinning up %v workers\n", workers) for i := 0; i < workers; i++ { go listStuff(&wg, i, inputChan, outputChan) } // 单独goroutine发送数据,避免主goroutine阻塞在发送操作上 go func() { for i := 0; i < 3000; i++ { inputChan <- i } // 发送完所有数据后关闭inputChan,告知worker没有更多数据 close(inputChan) }() // 等待所有worker完成后关闭outputChan go func() { wg.Wait() close(outputChan) }() // 从outputChan读取所有数据,直到通道关闭 for o := range outputChan { fmt.Println("reading from channel...", o) _output = append(_output, o) } fmt.Printf("+++ output len: %v\n", len(_output)) return _output, nil } func main() { List(5) }
修复说明
- 异步发送数据:将向
inputChan发送数据的逻辑放到单独goroutine中,避免主goroutine阻塞在发送操作,保证主goroutine可以持续从outputChan读取数据,缓解通道缓冲满的问题。 - 正确的通道关闭顺序:
- 数据发送完成后立即关闭
inputChan,让worker的for range循环在读取完所有数据后退出。 - 等待所有worker完成(
wg.Wait())后再关闭outputChan,确保所有数据都已转发到outputChan,主goroutine可以读取到全部数据。
- 数据发送完成后立即关闭
- 消除死锁:主goroutine持续读取
outputChan,避免worker因outputChan缓冲满而阻塞,worker就能持续从inputChan取数据,保证数据流动不中断。
内容的提问来源于stack exchange,提问作者salman b
相关产品推荐
相关产品推荐

