代码拉取完成,页面将自动刷新
package scheduler
import (
"gitee.com/h79/goutils/common"
"sync/atomic"
)
// 0:done, 1:working, -1:release
type StateCallFunc func(state int)
type ResultCallFunc func(state *JobResult)
type Pool struct {
jobs chan *Job
workers chan *jobWorker
stop chan bool
numWorkers int32 //总数
runWorkers int32 //正常运行
idleWorking int32 //空闲worker数
resultFn ResultCallFunc
stateFn StateCallFunc
running common.RunningCheck
}
func NewPool(numWorkers, jobQueueLen int) *Pool {
if numWorkers <= 0 || jobQueueLen <= 0 {
panic("length is zero")
}
pool := &Pool{
jobs: make(chan *Job, jobQueueLen),
workers: make(chan *jobWorker, numWorkers),
stop: make(chan bool),
numWorkers: int32(numWorkers),
runWorkers: int32(numWorkers),
idleWorking: int32(numWorkers),
resultFn: nil,
stateFn: nil,
}
pool.createWorkers(numWorkers)
pool.Run()
return pool
}
func (p *Pool) WithResultFunc(fn ResultCallFunc) *Pool {
p.resultFn = fn
return p
}
func (p *Pool) WithStateFunc(fn StateCallFunc) *Pool {
p.stateFn = fn
return p
}
func (p *Pool) AddJob(job *Job) {
if job == nil {
return
}
ww := atomic.LoadInt32(&p.runWorkers)
if ww <= p.numWorkers/2 {
// 中间有原因,线程池中的线程,由于某些原因,退出了,需要再启动
p.createWorkers(int(p.numWorkers - ww))
}
p.jobs <- job
}
func (p *Pool) Run() {
p.running.TryGoRunning(p.dispatch)
}
func (p *Pool) HasWorking() bool {
return atomic.LoadInt32(&p.idleWorking) != p.numWorkers
}
func (p *Pool) Release() {
p.stop <- true
<-p.stop
p.stateNotify(-1)
}
func (p *Pool) createWorkers(num int) {
for i := 0; i < num; i++ {
worker := newWorker(p)
worker.start()
}
}
func (p *Pool) stateNotify(state int) {
if p.stateFn != nil {
p.stateFn(state)
}
}
func (p *Pool) workerRunning() {
atomic.AddInt32(&p.runWorkers, 1)
}
func (p *Pool) workerQuit() {
atomic.AddInt32(&p.runWorkers, -1)
}
func (p *Pool) dispatch() {
for {
select {
case job := <-p.jobs:
worker := <-p.workers
atomic.AddInt32(&p.idleWorking, -1)
p.stateNotify(1)
worker.timeout = job.JTimeout
worker.job <- job
case stop := <-p.stop:
if stop {
for i := 0; i < cap(p.workers); i++ {
worker := <-p.workers
worker.stop <- true
<-worker.stop
}
atomic.StoreInt32(&p.idleWorking, p.numWorkers)
p.stop <- true
return
}
}
}
}
func (p *Pool) jobDone(result JobResult) {
atomic.AddInt32(&p.idleWorking, 1)
p.stateNotify(0)
if p.resultFn != nil {
p.resultFn(&result)
}
}
此处可能存在不合适展示的内容,页面不予展示。您可通过相关编辑功能自查并修改。
如您确认内容无涉及 不当用语 / 纯广告导流 / 暴力 / 低俗色情 / 侵权 / 盗版 / 虚假 / 无价值内容或违法国家有关法律法规的内容,可点击提交进行申诉,我们将尽快为您处理。