From 2a2f1347e89b2361600bb45f8e200251d4023f0e Mon Sep 17 00:00:00 2001 From: Johannes Findeisen Date: Sun, 27 Oct 2013 00:18:30 +0200 Subject: [PATCH] added some more intelligent logic to job scheduling. they are all scheduled in few seconds using a delay for each job instead of sleep --- bin/linspector | 12 +++++++++--- linspector/config/periods.py | 18 +++++++++++++----- linspector/linspector.py | 11 +++++++---- 3 files changed, 29 insertions(+), 12 deletions(-) diff --git a/bin/linspector b/bin/linspector index ca3b12c..5b3d534 100755 --- a/bin/linspector +++ b/bin/linspector @@ -73,9 +73,9 @@ def setup_logging(logfile="./log/linspector.log", logLevel=logging.ERROR, logfil if not path.exists(path.dirname(logfile)): os.makedirs(path.dirname(logfile)) - #logging.basicConfig(level=logging.WARNING) + logging.basicConfig(level=logging.WARNING) log = logging.getLogger(__name__) - log.setLevel(logging.DEBUG) + #log.setLevel(logging.DEBUG) consoleHandler = logging.StreamHandler() consoleHandler.setLevel(logLevel) @@ -105,7 +105,8 @@ def main(): configParser = FullConfigParser(log) linConf, core = configParser.parse_config(args.config) - scheduler = Scheduler() + scheduler = Scheduler({"apscheduler.threadpool.max_threads": 1000}) + print scheduler._threadpool scheduler.start() q = Queue.Queue() @@ -123,6 +124,11 @@ def main(): jsonrpc.daemon = True jsonrpc.start() + # Just for debug purposes: + #while True: + # print "Threads: " + '%d/%d' % (scheduler._threadpool.num_threads, scheduler._threadpool.max_threads) + # time.sleep(0.5) + LishFrontend(interface) log.debug("shutting down scheduler") diff --git a/linspector/config/periods.py b/linspector/config/periods.py index cb94577..24125a4 100644 --- a/linspector/config/periods.py +++ b/linspector/config/periods.py @@ -42,9 +42,13 @@ class IntervalPeriod(Period): self.start_date = start_date # when to first execute self.comment = comment # comment - def createJob(self, scheduler, jobInfo, func): + def createJob(self, scheduler, jobInfo, func, **kwargs): + start_date = self.start_date + if kwargs["start_date"]: + start_date = kwargs["start_date"] + return scheduler.add_interval_job(func, weeks=self.weeks, hours=self.hours, minutes=self.minutes, - seconds=self.seconds, start_date=self.start_date, args=[jobInfo]) + seconds=self.seconds, start_date=start_date, args=[jobInfo]) def __str__(self): ret = "IntervalPeriod(Name: " + self.name + ")" @@ -70,10 +74,14 @@ class CronPeriod(Period): ret = "CronPeriod(Name: " + self.name + ")" return ret - def createJob(self, scheduler, jobInfo, func): + def createJob(self, scheduler, jobInfo, func, **kwargs): + start_date = self.start_date + if kwargs["start_date"]: + start_date = kwargs["start_date"] + return scheduler.add_cron_job(func, year=self.year, month=self.month, day=self.day, week=self.week, day_of_week=self.day_of_week, hour=self.hour, minute=self.minute, - second=self.second, start_date=self.start_date, args=[jobInfo]) + second=self.second, start_date=start_date, args=[jobInfo]) class DatePeriod(Period): @@ -86,7 +94,7 @@ class DatePeriod(Period): ret = "DatePeriod(Name: " + self.name + ", " + str(self.date) + ")" return ret - def createJob(self, scheduler, jobInfo, func): + def createJob(self, scheduler, jobInfo, func, **kwargs): try: return scheduler.add_date_job(func=func, date=self.date, args=[jobInfo]) except Exception, e: diff --git a/linspector/linspector.py b/linspector/linspector.py index c74a2e2..7dccf50 100644 --- a/linspector/linspector.py +++ b/linspector/linspector.py @@ -20,8 +20,8 @@ You should have received a copy of the GNU Affero General Public License along with this program. If not, see . """ +import datetime import threading -import time from core.job import Job @@ -43,22 +43,25 @@ class Linspector(threading.Thread): jobInfo.handle_call() def run(self): + start_date = datetime.datetime.now() + time_delta = 0 jobs = [] for layout in self.linConf.get_enabled_layouts(): for hostgroup in layout.get_hostgroups(): for service in hostgroup.get_services(): for host in hostgroup.get_hosts(): for period in service.get_periods(): + time_delta += 2 + new_start_date = start_date + datetime.timedelta(seconds=time_delta) job = Job(service, host, hostgroup.get_members(), hostgroup.get_processors(), self.core, hostgroup) - schedulerJob = period.createJob(self.scheduler, job, handle_job) + schedulerJob = period.createJob(self.scheduler, job, handle_job, start_date=new_start_date) if schedulerJob is not None: job.set_job(schedulerJob) job.set_logger(self.log) jobs.append(job) - self.q.put(jobs) - time.sleep(3) \ No newline at end of file + self.q.put(jobs) \ No newline at end of file