带缓冲通道出现all goroutines are asleep - deadlock的原因排查
固定Goroutine数处理任务时的死锁问题分析
我希望创建固定数量(例如5个)的goroutine来处理数量可变的任务,实现代码及测试代码如下。测试用例「jobs 10 capacity 5」可正常运行,但「jobs 100 capacity 5」执行失败;将capacity设为50时可正常运行,设为30则仍失败。我原本认为缓冲通道满时会阻塞直到有空闲容量,通过信号量sem控制goroutine数量不超过设定的capacity,但为何会出现fatal error: all goroutines are asleep - deadlock!错误?
原实现代码
package main import ( "context" "fmt" "runtime" "time" ) type Job struct { id int result bool } func doWork(size int, capacity int) int { start := time.Now() jobs := make(chan *Job, capacity) results := make(chan *Job, capacity) sem := make(chan struct{}, capacity) go chanWorker(jobs, results, sem) for i := 0; i < size; i++ { jobs <- &Job{id: i} } close(jobs) successCount := 0 for i := 0; i < size; i++ { item := <-results if item.result { successCount++ } fmt.Printf("Job %d completed %v\n", item.id, item.result) } close(results) close(sem) fmt.Printf("Time taken to execute %d jobs with %d capacity = %v\n", size, capacity, time.Since(start)) return successCount } func chanWorker(jobs <-chan *Job, results chan<- *Job, sem chan struct{}) { for item := range jobs { it := item sem <- struct{}{} fmt.Printf("Job %d started\n", it.id) go func() { timeOutCtx, cancel := context.WithTimeout(context.Background(), 300*time.Millisecond) defer cancel() time.Sleep(time.Duration(it.id) * 100 * time.Millisecond) select { case <-timeOutCtx.Done(): fmt.Printf("Job %d timed out\n", it.id) it.result = false results <- it <-sem return default: fmt.Printf("Total number of routines %d\n", runtime.NumGoroutine()) it.result = true results <- it <-sem } }() } }
测试代码
package main import ( "testing" ) func Test_doWork(t *testing.T) { type args struct { size int capacity int } tests := []struct { name string args args want int }{ { name: "jobs 10 capacity 5", args: args{ size: 10, capacity: 5, }, want: 3, }, { name: "jobs 100 capacity 5", args: args{ size: 100, capacity: 5, }, want: 3, }, } for _, tt := range tests { t.Run(tt.name, func(t *testing.T) { if got := doWork(tt.args.size, tt.args.capacity); got < tt.want { t.Errorf("doWork() = %v, want %v", got, tt.want) } }) } }
死锁原因分析
阻塞链形成:
jobs通道缓冲为capacity(如5),当任务量为100时,主函数往jobs塞任务会因通道满而阻塞,等待chanWorker取走任务。chanWorker是单个goroutine,每次取任务后先执行sem <- struct{}{}获取信号量,当sem被占满(5个),chanWorker会卡在这一步,无法继续从jobs取任务,导致主函数的jobs <-操作持续阻塞。- 同时,处理任务的goroutine完成后会往
results通道塞结果,results缓冲同样为capacity,当通道满时,这些goroutine会卡在results <- it,无法执行<-sem释放信号量,进一步导致chanWorker无法继续处理任务,最终所有goroutine进入阻塞状态,触发死锁。
逻辑设计缺陷:
原代码试图用信号量控制并发,但chanWorker本身是单goroutine,会因信号量阻塞而中断任务消费链路;同时每个任务都启动新goroutine的设计,也不符合「固定数量goroutine处理任务」的需求。
修复方案
改为启动固定数量的worker goroutine,每个worker循环处理任务,无需额外信号量即可严格控制并发数:
package main import ( "context" "fmt" "runtime" "time" ) type Job struct { id int result bool } func doWork(size int, capacity int) int { start := time.Now() jobs := make(chan *Job, capacity) results := make(chan *Job, capacity) // 启动固定数量的worker goroutine for w := 0; w < capacity; w++ { go worker(jobs, results) } // 单独goroutine发送任务,避免阻塞主函数 go func() { for i := 0; i < size; i++ { jobs <- &Job{id: i} } close(jobs) }() successCount := 0 for i := 0; i < size; i++ { item := <-results if item.result { successCount++ } fmt.Printf("Job %d completed %v\n", item.id, item.result) } close(results) fmt.Printf("Time taken to execute %d jobs with %d capacity = %v\n", size, capacity, time.Since(start)) return successCount } func worker(jobs <-chan *Job, results chan<- *Job) { for it := range jobs { fmt.Printf("Job %d started\n", it.id) timeOutCtx, cancel := context.WithTimeout(context.Background(), 300*time.Millisecond) defer cancel() time.Sleep(time.Duration(it.id) * 100 * time.Millisecond) select { case <-timeOutCtx.Done(): fmt.Printf("Job %d timed out\n", it.id) it.result = false results <- it default: fmt.Printf("Total number of routines %d\n", runtime.NumGoroutine()) it.result = true results <- it } } }
修复说明
- 启动
capacity个固定worker,每个worker循环从jobs取任务,并发数严格控制为设定值。 - 任务发送逻辑放在单独goroutine中,避免主函数因
jobs通道满而阻塞。 - 主函数专注消费
results结果,确保results通道不会因积压而阻塞worker,打破原有的阻塞闭环。
内容的提问来源于stack exchange,提问作者ykesh
相关产品推荐
相关产品推荐

