forked from contribsys/faktory
/
scheduler.go
110 lines (95 loc) · 2.42 KB
/
scheduler.go
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
package server
import (
"encoding/json"
"sync"
"sync/atomic"
"time"
"github.com/contribsys/faktory"
"github.com/contribsys/faktory/storage"
"github.com/contribsys/faktory/util"
)
type scanner struct {
name string
adapter scannerAdapter
jobs int64
walltime int64
cycles int64
}
var (
defaultDelay = 1 * time.Second
)
func (s *scanner) Name() string {
return s.name
}
func (s *scanner) Execute() error {
count := int64(0)
start := time.Now()
elms, err := s.adapter.Prune(util.Nows())
if err != nil {
return err
}
if len(elms) > 0 {
util.Infof("%s processing %d jobs", s.name, len(elms))
}
for _, elm := range elms {
var job faktory.Job
err := json.Unmarshal(elm, &job)
if err != nil {
util.Error("Unable to unmarshal json", err, elm)
continue
}
err = s.adapter.Push(job.Queue, elm)
if err != nil {
util.Warnf("Error pushing job to '%s': %s", job.Queue, err.Error())
continue
}
count += 1
}
end := time.Now()
atomic.AddInt64(&s.cycles, 1)
atomic.AddInt64(&s.jobs, int64(count))
atomic.AddInt64(&s.walltime, end.Sub(start).Nanoseconds())
return nil
}
type rocksAdapter struct {
store storage.Store
ts storage.SortedSet
goodbye bool
}
func (ra *rocksAdapter) Push(name string, elm []byte) error {
if ra.goodbye {
// when the Dead elements come up for "scheduling", they have
// expired and are removed forever. Goodbye!
util.Debug("Removing dead element forever. Goodbye.")
return nil
}
que, err := ra.store.GetQueue(name)
if err != nil {
return err
}
return que.Push(elm)
}
func (ra *rocksAdapter) Prune(string) ([][]byte, error) {
return ra.ts.RemoveBefore(util.Nows())
}
func (ra *rocksAdapter) Size() int64 {
return ra.ts.Size()
}
type scannerAdapter interface {
Prune(string) ([][]byte, error)
Push(string, []byte) error
Size() int64
}
func (s *scanner) Stats() map[string]interface{} {
return map[string]interface{}{
"enqueued": s.jobs,
"cycles": s.cycles,
"size": s.adapter.Size(),
"wall_time_sec": (float64(s.walltime) / 1000000000),
}
}
func (s *Server) startScanners(waiter *sync.WaitGroup) {
s.taskRunner.AddTask(5, &scanner{name: "Scheduled", adapter: &rocksAdapter{s.store, s.store.Scheduled(), false}})
s.taskRunner.AddTask(5, &scanner{name: "Retries", adapter: &rocksAdapter{s.store, s.store.Retries(), false}})
s.taskRunner.AddTask(60, &scanner{name: "Dead", adapter: &rocksAdapter{s.store, s.store.Dead(), true}})
}