You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

带缓冲通道出现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)
            }
        })
    }
}

死锁原因分析

  1. 阻塞链形成:

    • 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进入阻塞状态,触发死锁。
  2. 逻辑设计缺陷:
    原代码试图用信号量控制并发,但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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.24 21:53:21