Go语言chan chan构造引发死锁问题:原因排查与修复求助
Go语言chan chan通道池死锁问题分析与修复
问题原因
- 通道池未回收worker通道:主协程从
pool取出worker通道发送任务后,没有将通道放回池里。当分发到第4个任务时,池里已经没有可用通道,主协程阻塞在workerChan := <-pool,直接触发死锁。 - Worker协程无法正常退出:Worker的
jobs通道始终没被关闭,for job := range jobs会一直阻塞等待新任务,导致defer wg.Done()永远无法执行,主协程的wg.Wait()会无限等待。 - WaitGroup使用逻辑错误:主协程给每个任务调用
wg.Add(1),但Worker的wg.Done()是在协程完全结束时才执行,而非每个任务完成后,这导致WaitGroup的计数和实际任务/协程的逻辑不匹配,进一步加剧阻塞。
修复方案
package main import ( "fmt" "sync" ) type Job struct { ID int } func worker(id int, jobs <-chan Job, pool chan<- chan Job, wg *sync.WaitGroup) { defer wg.Done() fmt.Printf("Worker %d starting\n", id) for job := range jobs { fmt.Printf("Worker %d processing job %d\n", id, job.ID) // 任务完成后把当前worker的通道放回池,供后续任务复用 pool <- jobs } fmt.Printf("Worker %d done\n", id) } func main() { numWorkers := 3 maxJobs := 10 var workerWg sync.WaitGroup // 创建带缓冲的worker通道池 pool := make(chan chan Job, numWorkers) // 给worker协程的WaitGroup计数 workerWg.Add(numWorkers) for i := 0; i < numWorkers; i++ { workerChan := make(chan Job) pool <- workerChan go worker(i, workerChan, pool, &workerWg) } // 用单独的WaitGroup跟踪所有任务完成状态 var jobWg sync.WaitGroup jobWg.Add(maxJobs) for i := 0; i < maxJobs; i++ { job := Job{ID: i} workerChan := <-pool // 异步发送任务,避免主协程被worker处理速度阻塞 go func(c chan Job, j Job) { defer jobWg.Done() c <- j }(workerChan, job) } // 先等所有任务处理完成 jobWg.Wait() // 关闭通道池,遍历关闭所有worker的任务通道,让worker退出循环 close(pool) for workerChan := range pool { close(workerChan) } // 等待所有worker协程完全退出 workerWg.Wait() fmt.Println("All jobs are processed") }
修复关键点
- 回收worker通道:Worker完成单个任务后,将自身的任务通道放回通道池,保证后续任务能获取到可用的worker通道。
- 拆分WaitGroup:用
workerWg等待所有worker协程退出,jobWg等待所有任务完成,避免计数逻辑混乱。 - 正确关闭通道:先等所有任务完成,再关闭通道池,遍历池中的worker通道并关闭,触发Worker的
for range循环退出,进而执行wg.Done()。 - 异步发送任务:主协程用匿名协程发送任务,避免因worker处理速度慢导致主协程阻塞。
内容的提问来源于stack exchange,提问作者knightcool
相关产品推荐
相关产品推荐

