Run start running queues
()
| 29 | |
| 30 | // Run start running queues |
| 31 | func (q *Queue) Run() { |
| 32 | if atomic.LoadUint32(&q.running) == 1 { |
| 33 | return |
| 34 | } |
| 35 | |
| 36 | atomic.StoreUint32(&q.running, 1) |
| 37 | for i := 0; i < q.maxWorkers; i++ { |
| 38 | q.workers[i] = newWorker(q.workerPool, q.wg) |
| 39 | q.workers[i].Start() |
| 40 | } |
| 41 | |
| 42 | go q.dispatcher() |
| 43 | } |
| 44 | |
| 45 | func (q *Queue) dispatcher() { |
| 46 | for job := range q.jobQueue { |