Go并发Worker Pool死锁求助:已出预期结果却触发死锁
Go Worker Pool死锁问题分析与修复
你的Worker Pool程序能正常输出结果但触发死锁,核心原因有三个:
任务通道未关闭:
addJobs向jobsCh写完所有任务后,没有关闭通道。worker协程里的for num := range jobsCh会一直阻塞等待新任务,永远无法执行wg2.Done(),导致后续流程卡住。未等待worker完成任务:
addWorkers启动worker后立刻调用wg.Done(),没有等待所有worker处理完任务。且worker全部完成后未关闭resultsCh,导致readResults的读取循环一直阻塞。结果通道未关闭:
readResults通过range读取resultsCh,如果通道不关闭,循环会一直等待新数据,无法退出执行wg.Done()。
修复后的代码
package main import ( "fmt" "sync" ) func addJobs(jobsCh chan<- int, wg *sync.WaitGroup) { for i := 1; i < 101; i++ { jobsCh <- i } close(jobsCh) // 写完所有任务后关闭通道,通知worker无新任务 wg.Done() } func worker(jobsCh <-chan int, resultsCh chan<- int, wg2 *sync.WaitGroup) { defer wg2.Done() // 用defer确保函数退出时必执行Done for num := range jobsCh { resultsCh <- num * num } } func addWorkers(jobsCh <-chan int, resultsCh chan<- int, wg *sync.WaitGroup) { var wg2 sync.WaitGroup for i := 0; i < 10; i++ { wg2.Add(1) go worker(jobsCh, resultsCh, &wg2) } wg2.Wait() // 等待所有worker完成任务 close(resultsCh) // 关闭结果通道,通知readResults停止读取 wg.Done() } func readResults(resultsCh <-chan int, wg *sync.WaitGroup) { for sq := range resultsCh { fmt.Printf("%v ", sq) } wg.Done() } func main() { var wg sync.WaitGroup jobsCh := make(chan int) resultsCh := make(chan int) wg.Add(1) go addJobs(jobsCh, &wg) wg.Add(1) go addWorkers(jobsCh, resultsCh, &wg) wg.Add(1) go readResults(resultsCh, &wg) wg.Wait() }
关键修改说明
addJobs末尾添加close(jobsCh),让worker知晓任务已全部发送,循环可以正常结束。addWorkers中先执行wg2.Wait()等待所有worker完成任务,再关闭resultsCh,确保所有结果写入后通知读取端停止。worker用defer wg2.Done(),保证无论函数如何退出都能标记任务完成,避免遗漏。
内容的提问来源于stack exchange,提问作者ankush981
相关产品推荐
相关产品推荐

