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

如何统计特定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阻塞。
  • 最小worker数保障:Judge函数中仅在当前worker数大于最小值时才会发送kill信号,确保worker数量不会低于初始设定的5个。

内容的提问来源于stack exchange,提问作者HHJ

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.16 08:35:22