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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.07 02:09:03