Go支持动态新增任务的作业队列如何优雅判断所有Worker空闲并退出?
问题解答
一、更优雅的实现方案
你原来跟踪空闲Worker的方案冗余且存在数据竞争隐患,推荐采用**待处理任务计数 + sync.WaitGroup**的模式,核心逻辑如下:
- 完全不需要跟踪空闲Worker状态,只需要记录当前所有未完成的任务总数(包含正在处理的、队列中等待的)
- 初始任务入队前调用
wg.Add(1)计数 - 每次Worker生成新任务前,先调用
wg.Add(1)再将任务加入队列,保证计数不会漏 - 每个Worker处理完单个任务后调用
wg.Done()减少计数 - 主协程只需要调用
wg.Wait(),等计数归零就说明所有任务(包括后续动态新增的)全部处理完成,直接退出即可
该方案也能解决你担心的channel死锁问题:只要保证发送任务前已经调用了wg.Add,就不会出现所有Worker都阻塞在发送的情况——因为总计数是准确的,只要还有未完成的任务,就一定有Worker正在执行任务,后续会从channel取走任务。
二、Worker数量过大性能暴跌的原因
你观察到的性能问题来自两个核心瓶颈:
- Goroutine调度开销过高:Go的goroutine虽然轻量,但调度依然有开销。你的CPU物理核心通常只有几到几十核,开30万goroutine会导致Go运行时频繁进行goroutine上下文切换,大量CPU时间消耗在调度而非实际任务处理上,属于典型的“过饱和调度”。
- 全局锁竞争极端严重:你所有操作队列、修改空闲计数的逻辑都抢同一把全局互斥锁,30万goroutine同时抢锁的场景下,90%以上的时间都消耗在锁等待上,实际处理任务的CPU占比极低。
三、优化后的实现代码
共享队列版本
package main import ( "math/rand" "sync" ) const NB_WORKERS = 3000 const NB_INITIAL_JOBS = 300 const JOB_COUNTING_TO = 10000000 var ( jobQueue []int mu sync.Mutex wg sync.WaitGroup ) func worker() { for { // 取任务 mu.Lock() if len(jobQueue) == 0 { mu.Unlock() return } job := jobQueue[0] jobQueue = jobQueue[1:] mu.Unlock() // 执行任务 for i := 0; i < job; i++ { } // 概率新增任务 if rand.Intn(5) == 0 { wg.Add(1) // 新增任务前先计数 mu.Lock() jobQueue = append(jobQueue, JOB_COUNTING_TO) mu.Unlock() } // 标记当前任务处理完成 wg.Done() } } func main() { // 初始化任务 for i := 0; i < NB_INITIAL_JOBS; i++ { wg.Add(1) jobQueue = append(jobQueue, JOB_COUNTING_TO) } // 启动固定数量Worker for i := 0; i < NB_WORKERS; i++ { go worker() } // 等待所有任务处理完成 wg.Wait() }
无锁Channel版本(更简洁)
不需要手动维护队列和锁,直接用带缓冲Channel实现任务队列,完全规避死锁问题:
package main import ( "math/rand" "sync" ) const NB_WORKERS = 3000 const NB_INITIAL_JOBS = 300 const JOB_COUNTING_TO = 10000000 var ( jobChan = make(chan int, 1024) wg sync.WaitGroup ) func worker() { for job := range jobChan { // 执行任务 for i := 0; i < job; i++ { } // 概率新增任务 if rand.Intn(5) == 0 { wg.Add(1) jobChan <- JOB_COUNTING_TO } wg.Done() } } func main() { // 启动Worker for i := 0; i < NB_WORKERS; i++ { go worker() } // 初始化任务 for i := 0; i < NB_INITIAL_JOBS; i++ { wg.Add(1) jobChan <- JOB_COUNTING_TO } // 等待所有任务完成 wg.Wait() close(jobChan) // 关闭channel让所有Worker退出 }
额外优化建议
- Worker数量不要设置为远大于CPU核心数,通常设置为
runtime.NumCPU() * 2到runtime.NumCPU() * 4即可,避免调度开销过大 - 如果任务量极大,可以考虑用分段队列、局部锁的方式减少全局锁的竞争,进一步提升性能
内容的提问来源于stack exchange,提问作者Cosmo Sterin
相关产品推荐
相关产品推荐

