恒美微站
首页
关于我们
建站服务
主题模板
案例展示
资讯中心
联系我们
go 学习 - prometheus 指标协程池监控
首页
资讯中心
/
go 学习 - prometheus 指标协程池监控
go 学习 - prometheus 指标协程池监控
发布时间:2026/10/9 2:07:57
文章目录1. 协程池2. 监控当前协程池运行指标3. 小结1. 协程池go 里面都是用的协程由于协程创建比较轻量所以一般在业务里面要用就直接开一个协程很少有用协程池的但是有一些项目需要处理高并发问题QPS 能打到 10000 以上这种情况下就会考虑用协程池。go 里面用的比较多的就是ants需要在 go.mod 里面引入 ants 包github.com/panjf2000/ants/v2 v2.12.0。接下来就是初始化初始化比较简单。funcNew(cfg Config)(*Pool,error){p,err:ants.NewPool(cfg.WorkerCount,ants.WithMaxBlockingTasks(cfg.MaxBlockingTasks),ants.WithNonblocking(cfg.Nonblocking),ants.WithPreAlloc(cfg.PreAlloc),)iferr!nil{returnnil,err}returnPool{pool:p},nil}我们主要看下里面几个参数原理就先不看了。WorkerCount协程池子的大小也就是协程数相当于线程池的最大线程数。MaxBlockingTasks最多允许多少个调用方阻塞在pool.Submit()上等待0 默认表示不限制。Nonblocking是否非阻塞提交如果开启就是Submit()不等待池满了立刻返回ErrPoolOverload如果不开启就是池子满了就阻塞等待当Nonblocking false的时候上面的 MaxBlockingTasks 才有用。PanicHandlerworker 执行任务发生 panic 时的处理函数可以统一打印日志上报指标等。Logger自定义日志器如果不设置就用标准库 log通常用来接入自己的日志框架。DisablePurge是否禁用空闲 worker 清理如果是 true 就不清理worker 常驻如果是 false 就允许根据 ExpiryDuration 回收空闲 worker。ExpiryDuration空闲 worker 的过期时间协程池会有一个清理协程周期性扫描 worker如果某个 worker 超过这个时间没被使用就会被回收。上面是几个核心参数用法也很简单调用 Submit 添加任务。func(p*Pool)Submit(taskfunc())error{ifp.IsClosed(){returnErrPoolClosed}w,err:p.retrieveWorker()ifw!nil{w.inputFunc(task)}returnerr}可以看到任务就是一个 func() 无参无返回值函数。2. 监控当前协程池运行指标ants 协程池提供下面几个方法来监控协程池运行时候的协程数监控协程数可以看到这个协程池协程的利用率然后监控阻塞数就能看到是不是协程参数配置有问题导致一些协程阻塞住这下就要考虑调大协程数了。packagepoolimport(github.com/panjf2000/ants/v2)typeTaskfunc()typeConfigstruct{// WorkerCount 表示协程池最大并发 worker 数WorkerCountint// MaxBlockingTasks 表示最多允许多少个提交方阻塞等待MaxBlockingTasksint// Nonblocking 为 true 时池满后直接返回错误不阻塞等待Nonblockingbool// PreAlloc 为 true 时预分配内部 worker 队列PreAllocbool}typePoolstruct{pool*ants.Pool}funcNew(cfg Config)(*Pool,error){p,err:ants.NewPool(cfg.WorkerCount,ants.WithMaxBlockingTasks(cfg.MaxBlockingTasks),ants.WithNonblocking(cfg.Nonblocking),ants.WithPreAlloc(cfg.PreAlloc),)iferr!nil{returnnil,err}returnPool{pool:p},nil}func(p*Pool)Submit(task Task)error{returnp.pool.Submit(task)}// Running 返回正在运行的 worker 协程数func(p*Pool)Running()int{returnp.pool.Running()}// Free 返回当前空闲的 worker 数func(p*Pool)Free()int{returnp.pool.Free()}// Cap 返回当前协程池容量上限func(p*Pool)Cap()int{returnp.pool.Cap()}// Waiting 有多少阻塞等待提交的任务, 类似阻塞队列func(p*Pool)Waiting()int{returnp.pool.Waiting()}func(p*Pool)Close(){p.pool.Release()}然后我们可以启动一个定时任务每秒定时上报里面的指标下面就是 metrics 文件的内容。packagemetricsimport(fmtnet/httptimeexample.com/prometheus-2/poolgithub.com/prometheus/client_golang/prometheusgithub.com/prometheus/client_golang/prometheus/promhttp)typeRegistrystruct{registry*prometheus.Registry running prometheus.Gauge free prometheus.Gauge capacity prometheus.Gauge waiting prometheus.Gauge taskDuration prometheus.Histogram}funcNewRegistry()*Registry{reg:prometheus.NewRegistry()r:Registry{registry:reg,running:prometheus.NewGauge(prometheus.GaugeOpts{Name:pool_running_workers,Help:Current running workers.,}),free:prometheus.NewGauge(prometheus.GaugeOpts{Name:pool_free_workers,Help:Current free workers.,}),capacity:prometheus.NewGauge(prometheus.GaugeOpts{Name:pool_capacity,Help:Current pool capacity.,}),waiting:prometheus.NewGauge(prometheus.GaugeOpts{Name:pool_waiting_tasks,Help:Current pool waiting tasks.,}),taskDuration:prometheus.NewHistogram(prometheus.HistogramOpts{Name:pool_task_duration_seconds,Help:Duration of tasks executed inside workerPool.Submit.,Buckets:prometheus.DefBuckets,}),}reg.MustRegister(r.running,r.free,r.capacity,r.waiting,r.taskDuration)returnr}func(r*Registry)Handler()http.Handler{returnpromhttp.HandlerFor(r.registry,promhttp.HandlerOpts{})}func(r*Registry)WrapTaskDuration(taskfunc())func(){returnfunc(){start:time.Now()deferfunc(){r.taskDuration.Observe(time.Since(start).Seconds())}()task()}}func(r*Registry)RunTicker(pool*pool.Pool){ticker:time.NewTicker(1*time.Second)deferticker.Stop()for{select{case-ticker.C:fmt.Println(Now() 上报指标)r.running.Set(float64(pool.Running()))r.free.Set(float64(pool.Free()))r.capacity.Set(float64(pool.Cap()))r.waiting.Set(float64(pool.Waiting()))}}}funcNow()string{returntime.Now().Format(2006-01-02 15:04:05.000)}最后就是 main 方法启动main 方法中我们启动一个定时任务每秒上报协程池的几个指标然后启动一个定时任务每秒添加 50 个任务到协程池里面。packagemainimport(fmtlogmath/randnet/httptimeexample.com/prometheus-2/metricsexample.com/prometheus-2/pool)constsubmitBatchSize50funcmain(){reg:metrics.NewRegistry()workerPool,err:pool.New(pool.Config{WorkerCount:50,MaxBlockingTasks:128,Nonblocking:true,PreAlloc:false,})iferr!nil{log.Fatalf(create worker pool failed: %v,err)}goreg.RunTicker(workerPool)gostartSubmitTicker(workerPool,reg)deferworkerPool.Close()http.HandleFunc(/,func(w http.ResponseWriter,req*http.Request){_,_fmt.Fprintln(w,hello prometheus-2)})http.HandleFunc(/healthz,func(w http.ResponseWriter,req*http.Request){_,_w.Write([]byte(ok))})http.Handle(/metrics,reg.Handler())addr::8088log.Printf(prometheus-2 listening on %s,addr)log.Printf(metrics endpoint: http://127.0.0.1%s/metrics,addr)log.Fatal(http.ListenAndServe(addr,nil))}funcstartSubmitTicker(workerPool*pool.Pool,reg*metrics.Registry){ticker:time.NewTicker(1000*time.Millisecond)deferticker.Stop()forrangeticker.C{fori:0;isubmitBatchSize;i{taskID:i1task:reg.WrapTaskDuration(func(){cost:1rand.Intn(100)time.Sleep(time.Duration(cost)*time.Millisecond)log.Printf(task %d finished, cost%dms,taskID,cost)})iferr:workerPool.Submit(task);err!nil{log.Printf(submit task %d failed: %v,taskID,err)break}}}}funcinit(){rand.Seed(time.Now().UnixNano())}最后来看下输出结果打开 prometheus 输入 promql 查看监控指标。当前协程池的当前运行协程数。当前协程池里面的空闲协程数。当前协程池的协程容量大小。阻塞等待提交的任务数。请求平均耗时(rate(pool_task_duration_seconds_sum[1m]) / rate(pool_task_duration_seconds_count[1m])) * 1000。请求 p95 耗时** histogram_quantile(0.95, rate(pool_task_duration_seconds_bucket[1m])) * 1000**。3. 小结这篇文章简单学习下 prometheus 协程池监控的时候核心就是希望通过监控协程池发现运行过程中的性能问题。如有错误欢迎指出