如何统计特定goroutine数量?实现自适应goroutine工作池
动态增减Goroutine工作池实现方案
要实现不依赖全局变量的动态worker池,核心是安全跟踪活跃worker数量,同时保证worker数量不低于设定的最小值(此处为5)。以下是具体实现:
修改后的完整代码
package main import ( "fmt" "sync/atomic" "time" ) func Adaptive_Worker_Pool(value_input chan int) { kill_sig := make(chan bool) make_sig := make(chan bool) var workerCount int32 // 用原子变量跟踪活跃worker数量,避免全局变量 // 初始化5个worker initCount := 5 for i := 0; i < initCount; i++ { atomic.AddInt32(&workerCount, 1) go Do(kill_sig, value_input, &workerCount) } go Make_Routine(make_sig, kill_sig, value_input, &workerCount) go Judge(kill_sig, make_sig, value_input, &workerCount, initCount) } func Make_Routine(make_sig chan bool, kill_sig chan bool, value_input chan int, workerCount *int32) { for { <-make_sig atomic.AddInt32(workerCount, 1) go Do(kill_sig, value_input, workerCount) } } func Do(kill_sig chan bool, value_input chan int, workerCount *int32) { defer atomic.AddInt32(workerCount, -1) // 退出时自动减少计数 outer: for { select { case value := <-value_input: fmt.Println(value) case <-kill_sig: break outer } } } func Judge(kill_sig chan bool, make_sig chan bool, value_input chan int, workerCount *int32, minWorker int) { for { time.Sleep(time.Millisecond * 500) taskCount := len(value_input) currentWorkers := atomic.LoadInt32(workerCount) if taskCount > 5 { // 任务积压,新增worker make_sig <- true } else { // 任务不足,且当前worker数大于最小值时,减少worker if currentWorkers > int32(minWorker) { // 用select避免发送kill信号时阻塞(比如所有worker都在处理任务) select { case kill_sig <- true: default: } } } } } func main() { value_input := make(chan int, 10) Adaptive_Worker_Pool(value_input) a := 0 for { value_input <- a a++ } }
关键实现细节
- 原子变量计数:在
Adaptive_Worker_Pool内部定义workerCount原子整型,通过闭包传递给各goroutine,避免全局变量。原子操作保证多goroutine环境下计数的准确性和安全性。 - worker生命周期绑定计数:每个worker启动时调用
atomic.AddInt32增加计数,退出时通过defer自动减少计数,确保计数与worker实际数量严格同步。 - 动态调整逻辑:
- 当任务队列长度超过5时,发送
make_sig信号新增worker; - 当任务队列长度不足5,且当前worker数大于最小值(5)时,发送
kill_sig信号减少worker。发送kill信号时加入select+default,避免所有worker都在处理任务时无法接收信号导致Judge阻塞。
- 当任务队列长度超过5时,发送
- 最小worker数保障:Judge函数中仅在当前worker数大于最小值时才会发送kill信号,确保worker数量不会低于初始设定的5个。
内容的提问来源于stack exchange,提问作者HHJ
相关产品推荐
相关产品推荐

