1 Star 0 Fork 0

h79/goutils

加入 Gitee
与超过 1200万 开发者一起发现、参与优秀开源项目,私有仓库也完全免费 :)
免费加入
文件
克隆/下载
pool.go 3.43 KB
一键复制 编辑 原始数据 按行查看 历史
huqiuyun 提交于 2024-08-25 23:49 . scheduler job for system quit
package scheduler
import (
"gitee.com/h79/goutils/common/logger"
"gitee.com/h79/goutils/common/system"
"sync/atomic"
"time"
)
type IdleHandler interface {
Idle(num int32)
}
// IdleFunc idle working number
type IdleFunc func(num int32)
func (fn IdleFunc) Idle(num int32) {
fn(num)
}
type Pool struct {
jobs chan *Job
workers chan *jobWorker
stop chan bool
delay *Delay
tm *time.Ticker
numWorkers int32 //总数
runWorkers int32 //正常运行
idleWorking int32 //空闲worker数
idleHandler IdleHandler
running system.RunningCheck
}
func NewPool(numWorkers, jobQueueLen int) *Pool {
if numWorkers <= 0 || jobQueueLen <= 0 {
panic("length is zero")
}
var num = int32(numWorkers)
pool := &Pool{
jobs: make(chan *Job, jobQueueLen),
workers: make(chan *jobWorker, numWorkers),
stop: make(chan bool),
delay: NewDelay(),
tm: time.NewTicker(time.Millisecond * 10),
numWorkers: num,
runWorkers: num,
idleWorking: num,
}
pool.createWorkers(numWorkers)
pool.Run()
return pool
}
func (p *Pool) SetIdleHandler(fn IdleHandler) *Pool {
p.idleHandler = fn
return p
}
func (p *Pool) WithIdleFunc(fn func(num int32)) *Pool {
return p.SetIdleHandler(IdleFunc(fn))
}
func (p *Pool) AddJob(job *Job) {
if job == nil {
return
}
if system.IsQuit() {
logger.W("Task", "application is quited,not can add for jobId= %s", job.GetId())
return
}
ww := atomic.LoadInt32(&p.runWorkers)
if ww <= p.numWorkers/2 {
// 中间有原因,线程池中的线程,由于某些原因,退出了,需要再启动
p.createWorkers(int(p.numWorkers - ww))
}
p.Run()
p.jobs <- job
}
func (p *Pool) Run() {
p.running.GoRunning(p.dispatch)
}
func (p *Pool) HasWorking() bool {
return atomic.LoadInt32(&p.idleWorking) != p.numWorkers
}
func (p *Pool) Release() {
system.Stop(time.Second, p.stop)
p.idleHandle(p.numWorkers)
}
func (p *Pool) createWorkers(num int) {
for i := 0; i < num; i++ {
worker := newWorker(p)
worker.start()
}
}
func (p *Pool) idleHandle(idleWork int32) {
if p.idleHandler != nil {
p.idleHandler.Idle(idleWork)
}
}
func (p *Pool) workerRunning() {
atomic.AddInt32(&p.runWorkers, 1)
}
func (p *Pool) workerQuit() {
atomic.AddInt32(&p.runWorkers, -1)
}
func (p *Pool) workerIdle() {
p.idleHandle(atomic.AddInt32(&p.idleWorking, 1))
}
func (p *Pool) dispatch() {
defer p.tm.Stop()
defer p.idleHandle(p.numWorkers)
for {
select {
case <-p.stop:
p.stopWorks(1)
p.stop <- true
return
case <-system.Closed():
p.stopWorks(2)
return
case job := <-p.jobs:
if job.Delay > 0 {
p.delayJob(job)
} else {
p.executeJob(job)
}
case <-p.tm.C:
p.checkDelayJob()
}
}
}
func (p *Pool) stopWorks(s int) {
for i := 0; i < len(p.workers); i++ {
worker := <-p.workers
worker.stop <- s
<-worker.stop
}
atomic.StoreInt32(&p.idleWorking, p.numWorkers)
// delay
opts := With(s)
p.delay.Exec(opts...)
for i := 0; i < len(p.jobs); i++ {
job, ok := <-p.jobs
if !ok {
return
}
_, _ = job.Execute(opts...)
}
}
func (p *Pool) idleWork(w *jobWorker) {
if system.IsQuit() {
return
}
p.workers <- w
}
func (p *Pool) executeJob(job *Job) {
if system.IsQuit() {
_, _ = job.Execute(WithQuited())
return
}
worker := <-p.workers
p.idleHandle(atomic.AddInt32(&p.idleWorking, -1))
worker.job <- job
}
func (p *Pool) delayJob(job *Job) {
p.delay.Add(job)
}
func (p *Pool) checkDelayJob() {
p.delay.Check(p)
}
Loading...
马建仓 AI 助手
尝试更多
代码解读
代码找茬
代码优化
Go
1
https://gitee.com/h79/goutils.git
git@gitee.com:h79/goutils.git
h79
goutils
goutils
v1.22.7

搜索帮助