多goroutine读取同一channel异常:仅各读一次即停止
问题:多Goroutine从同一Channel读取时提前停止的原因
我启动了两个worker goroutine从同一个channel读取值,但每个worker只读取一个值就停止了。原本期望它们会持续读取直到生产者关闭channel,但现在生产者还没关闭,发送第三个值时就阻塞了,第三个值也没被任何worker读取。
运行输出
new worker new worker waiting sending 0 sending 1 sending 2 running func 1 sending value out 1 running func 0 sending value out 0
相关代码
package main import ( "fmt" "sync" ) func workerPool(done <-chan bool, in <-chan int, numberOfWorkers int, fn func(int) int) chan int { out := make(chan int) var wg sync.WaitGroup for i := 0; i < numberOfWorkers; i++ { fmt.Println("new worker") wg.Add(1) // 启动worker从in读取数据,处理后发送到out go func() { defer wg.Done() for { select { case <-done: fmt.Println("recieved done signal") return case data, ok := <-in: if !ok { fmt.Println("no more items") return } fmt.Println("running func", data) value := fn(data) fmt.Println("sending value out", value) out <- value } } }() } fmt.Println("waiting") wg.Wait() fmt.Println("done waiting") close(out) return out } func main() { done := make(chan bool) defer close(done) in := make(chan int) go func() { for i := 0; i < 10; i++ { fmt.Println("sending", i) in <- i } close(in) }() out := workerPool(done, in, 2, func(i int) int { return i }) for { select { case o, ok := <-out: if !ok { continue } fmt.Println("output", o) case <-done: return default: } } }
问题原因与解决方案
核心问题:双重阻塞导致死锁
Worker阻塞在结果发送:
out是无缓冲channel,worker处理完数据执行out <- value时,需要等待接收方(main函数)读取数据才能完成发送。但workerPool启动goroutine后立刻调用wg.Wait(),导致main函数无法进入读取out的循环,worker永远阻塞在out <- value,无法继续读取in的后续值。wg.Wait()位置错误:workerPool函数内直接调用wg.Wait(),但此时所有worker都阻塞在发送操作上,永远无法执行wg.Done(),导致wg.Wait()无限阻塞,workerPool无法返回,main函数彻底卡死在调用workerPool的步骤。
修复方案
方案1:调整等待逻辑,异步等待worker完成
修改workerPool,不要在函数主线程调用wg.Wait(),而是启动单独goroutine等待,让workerPool立刻返回out,让main函数可以开始读取结果:
func workerPool(done <-chan bool, in <-chan int, numberOfWorkers int, fn func(int) int) chan int { out := make(chan int) var wg sync.WaitGroup for i := 0; i < numberOfWorkers; i++ { fmt.Println("new worker") wg.Add(1) go func() { defer wg.Done() for { select { case <-done: fmt.Println("recieved done signal") return case data, ok := <-in: if !ok { fmt.Println("no more items") return } fmt.Println("running func", data) value := fn(data) fmt.Println("sending value out", value) out <- value } } }() } // 异步等待所有worker完成后关闭out go func() { wg.Wait() fmt.Println("done waiting") close(out) }() fmt.Println("worker pool started") return out }
方案2:简化main函数的结果读取逻辑
main函数里的读取循环可以直接遍历out,当out被关闭时循环会自动退出,无需额外的select判断:
func main() { done := make(chan bool) defer close(done) in := make(chan int) go func() { for i := 0; i < 10; i++ { fmt.Println("sending", i) in <- i } close(in) }() out := workerPool(done, in, 2, func(i int) int { return i }) // 直接遍历out,直到channel关闭 for o := range out { fmt.Println("output", o) } }
内容的提问来源于stack exchange,提问作者Jonathan Kittell
相关产品推荐
相关产品推荐

