From 4c5461d4ba7bc5bfcd1a9e41434c18eda1024611 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Torkel=20=C3=96degaard?= Date: Mon, 5 Sep 2016 14:26:08 +0200 Subject: [PATCH] feat(alerting): alerting scheduling distribution, only distibutes it on seconds for now, not sub second distribution, #5854 --- pkg/services/alerting/models.go | 9 +++++---- pkg/services/alerting/scheduler.go | 29 ++++++++++++++++++++++++----- 2 files changed, 29 insertions(+), 9 deletions(-) diff --git a/pkg/services/alerting/models.go b/pkg/services/alerting/models.go index e11e1e1aaaf..39e7eb83166 100644 --- a/pkg/services/alerting/models.go +++ b/pkg/services/alerting/models.go @@ -1,10 +1,11 @@ package alerting type Job struct { - Offset int64 - Delay bool - Running bool - Rule *Rule + Offset int64 + OffsetWait bool + Delay bool + Running bool + Rule *Rule } type ResultLogEntry struct { diff --git a/pkg/services/alerting/scheduler.go b/pkg/services/alerting/scheduler.go index ffac7ddb659..6701ad6eb6c 100644 --- a/pkg/services/alerting/scheduler.go +++ b/pkg/services/alerting/scheduler.go @@ -1,6 +1,7 @@ package alerting import ( + "math" "time" "github.com/grafana/grafana/pkg/log" @@ -34,8 +35,8 @@ func (s *SchedulerImpl) Update(rules []*Rule) { } job.Rule = rule - job.Offset = int64(i) - + job.Offset = ((rule.Frequency * 1000) / int64(len(rules))) * int64(i) + job.Offset = int64(math.Floor(float64(job.Offset) / 1000)) jobs[rule.Id] = job } @@ -46,9 +47,27 @@ func (s *SchedulerImpl) Tick(tickTime time.Time, execQueue chan *Job) { now := tickTime.Unix() for _, job := range s.jobs { - if now%job.Rule.Frequency == 0 && job.Running == false { - s.log.Debug("Scheduler: Putting job on to exec queue", "name", job.Rule.Name) - execQueue <- job + if job.Running { + continue + } + + if job.OffsetWait && now%job.Offset == 0 { + job.OffsetWait = false + s.enque(job, execQueue) + continue + } + + if now%job.Rule.Frequency == 0 { + if job.Offset > 0 { + job.OffsetWait = true + } else { + s.enque(job, execQueue) + } } } } + +func (s *SchedulerImpl) enque(job *Job, execQueue chan *Job) { + s.log.Debug("Scheduler: Putting job on to exec queue", "name", job.Rule.Name, "id", job.Rule.Id) + execQueue <- job +}