Updated to version 3 of the apscheduler library. multiprocessing should work now but some objects are not serializable so it does not. i will fix that some day.

This commit is contained in:
Johannes Findeisen 2015-06-23 01:25:15 +02:00
commit 40cb6c8d35
8 changed files with 41 additions and 25 deletions

View file

@ -1,6 +1,7 @@
#!/usr/bin/python2.7 -tt #!/usr/bin/python2.7 -tt
""" """
Copyright (c) 2014-2015 by Johannes Findeisen
Copyright (c) 2011-2013 by Johannes Findeisen and Rafael Timmerberg Copyright (c) 2011-2013 by Johannes Findeisen and Rafael Timmerberg
This file is part of Linspector (http://linspector.org). This file is part of Linspector (http://linspector.org).
@ -20,9 +21,10 @@ along with this program. If not, see <http://www.gnu.org/licenses/>.
""" """
__version__ = "0.15.1/AMNESIA" __version__ = "0.16/AMNESIA"
import argparse import argparse
import copy
import datetime import datetime
import logging import logging
import logging.handlers import logging.handlers
@ -33,10 +35,14 @@ import sys
from linspector.config.parser import FullConfigParser from linspector.config.parser import FullConfigParser
from linspector.core.interface import LinspectorInterface from linspector.core.interface import LinspectorInterface
from linspector.core.job import LinspectorJob from linspector.core.job import LinspectorJob
from linspector.core.scheduler import LinspectorScheduler #from linspector.core.scheduler import LinspectorScheduler
from linspector.frontends.lish import LishFrontend from linspector.frontends.lish import LishFrontend
from linspector.tasks.task import TaskExecutor from linspector.tasks.task import TaskExecutor
from apscheduler.schedulers.background import BackgroundScheduler
from apscheduler.jobstores.memory import MemoryJobStore
from apscheduler.executors.pool import ThreadPoolExecutor, ProcessPoolExecutor
logger = logging.getLogger(__name__) logger = logging.getLogger(__name__)
@ -67,8 +73,8 @@ def parse_args():
parser.add_argument("-t", "--threads", default=3500, type=int, parser.add_argument("-t", "--threads", default=3500, type=int,
help="maximum number of scheduler threads (default: 3500)") help="maximum number of scheduler threads (default: 3500)")
parser.add_argument("-k", "--corethreads", default=0, type=int, parser.add_argument("-p", "--processes", default=0, type=int,
help="number of scheduler core threads (default: 0)") help="number of scheduler processes (default: 0)")
parser.add_argument("-x", "--delay", default=1.135791, type=float, parser.add_argument("-x", "--delay", default=1.135791, type=float,
help="seconds delay between scheduled jobs (default: 1.135791)") help="seconds delay between scheduled jobs (default: 1.135791)")
@ -115,9 +121,18 @@ def main():
print("Configuration error: " + str(msg) + ". Exiting now.") print("Configuration error: " + str(msg) + ". Exiting now.")
logger.error(msg) logger.error(msg)
exit() exit()
scheduler = LinspectorScheduler({"apscheduler.threadpool.core_threads": args.corethreads, jobstores = {
"apscheduler.threadpool.max_threads": args.threads}) 'memory': MemoryJobStore()
}
executors = {
'default': ThreadPoolExecutor(args.threads),
'processpool': ProcessPoolExecutor(args.processes)
}
job_defaults = {
'max_instances': 10000
}
scheduler = BackgroundScheduler(jobstores=jobstores, executors=executors, job_defaults=job_defaults)
TaskExecutor.Instance() TaskExecutor.Instance()
@ -193,4 +208,4 @@ def main():
if __name__ == "__main__": if __name__ == "__main__":
main() main()

View file

@ -12,7 +12,7 @@
} }
}, },
"periods": { "periods": {
"fast": { "seconds": 60 }, "fast": { "seconds": 15 },
"medium": { "seconds": 120 }, "medium": { "seconds": 120 },
"slow": { "seconds": 480 } "slow": { "seconds": 480 }
}, },

View file

@ -52,8 +52,8 @@ class IntervalPeriod(Period):
if "start_date" in kwargs: if "start_date" in kwargs:
start_date = kwargs["start_date"] start_date = kwargs["start_date"]
return scheduler.add_interval_job(func, weeks=self.weeks, hours=self.hours, minutes=self.minutes, return scheduler.add_job(func, trigger="interval", weeks=self.weeks, hours=self.hours, minutes=self.minutes,
seconds=self.seconds, start_date=start_date, args=[jobInfo]) seconds=self.seconds, start_date=start_date, args=[jobInfo], timezone="CET")
def __str__(self): def __str__(self):
ret = "IntervalPeriod(Name: " + self.name + ")" ret = "IntervalPeriod(Name: " + self.name + ")"
@ -85,7 +85,7 @@ class CronPeriod(Period):
if "start_date" in kwargs: if "start_date" in kwargs:
start_date = 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, return scheduler.add_job(func, trigger="cron", 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, day_of_week=self.day_of_week, hour=self.hour, minute=self.minute,
second=self.second, start_date=start_date, args=[jobInfo]) second=self.second, start_date=start_date, args=[jobInfo])
@ -108,6 +108,6 @@ class DatePeriod(Period):
if date < earliest: if date < earliest:
self.date = earliest self.date = earliest
return scheduler.add_date_job(func=func, date=self.date, args=[jobInfo]) return scheduler.add_job(func=func, trigger="date", date=self.date, args=[jobInfo])
except Exception, e: except Exception, e:
logger.error("exception while creating job out of DatePeriod!\n%s" % e) logger.error("exception while creating job out of DatePeriod!\n%s" % e)

View file

@ -63,7 +63,7 @@ class LinspectorInterface(object):
d["Members"] = str([member.name for member in job.members]) d["Members"] = str([member.name for member in job.members])
d["Period"] = str(job.scheduler_job.trigger) d["Period"] = str(job.scheduler_job.trigger)
d["Next run"] = str(job.scheduler_job.next_run_time) d["Next run"] = str(job.scheduler_job.next_run_time)
d["Runs"] = str(job.scheduler_job.runs) d["Runs"] = str(job.job_information.job_overall_fails)
d["Enabled"] = str(job.enabled) d["Enabled"] = str(job.enabled)
d["Threshold"] = str(job.service.get_threshold()) d["Threshold"] = str(job.service.get_threshold())
d["Fails"] = str(job.job_threshold) d["Fails"] = str(job.job_threshold)

View file

@ -49,6 +49,7 @@ class LinspectorJob:
self.enabled = True self.enabled = True
self.scheduler_job = None self.scheduler_job = None
self.job_id = self.hex_string() self.job_id = self.hex_string()
""" """
NONE job was not executed NONE job was not executed
OK when everything is fine OK when everything is fine

View file

@ -1,4 +1,5 @@
""" """
Copyright (c) 2014 by Johannes Findeisen
Copyright (c) 2011-2013 by Johannes Findeisen and Rafael Timmerberg Copyright (c) 2011-2013 by Johannes Findeisen and Rafael Timmerberg
This file is part of Linspector (http://linspector.org). This file is part of Linspector (http://linspector.org).
@ -19,11 +20,12 @@ along with this program. If not, see <http://www.gnu.org/licenses/>.
from logging import getLogger from logging import getLogger
from apscheduler.scheduler import Scheduler from apscheduler.schedulers.background import BackgroundScheduler
from apscheduler.executors.pool import ThreadPoolExecutor, ProcessPoolExecutor
logger = getLogger(__name__) logger = getLogger(__name__)
class LinspectorScheduler(Scheduler): class LinspectorScheduler(Scheduler):
def test(self): def test(self):
pass pass

View file

@ -220,7 +220,6 @@ class NewLish(CmdWrapper):
count to get a job count''') count to get a job count''')
class Job(CmdWrapper): class Job(CmdWrapper):
def __init__(self, interface): def __init__(self, interface):
@ -478,8 +477,8 @@ Usage:
print GREEN + "Hostname" + END + ":\t" + socket.gethostname() print GREEN + "Hostname" + END + ":\t" + socket.gethostname()
# TODO: print instance uptime # TODO: print instance uptime
print GREEN + "Job Count" + END + ":\t" + str(self.interface.get_job_count()) print GREEN + "Job Count" + END + ":\t" + str(self.interface.get_job_count())
thread_info = self.interface.get_thread_count() ###thread_info = self.interface.get_thread_count()
print GREEN + "Threads" + END + ":\t" + str(thread_info["Num Threads"]) + "/" + str(thread_info["Max Threads"]) ###print GREEN + "Threads" + END + ":\t" + str(thread_info["Num Threads"]) + "/" + str(thread_info["Max Threads"])
print GREEN + "Core Version" + END + ":\t" + self.interface.get_version() print GREEN + "Core Version" + END + ":\t" + self.interface.get_version()
print GREEN + "Lish Version" + END + ":\t" + __version__ print GREEN + "Lish Version" + END + ":\t" + __version__
@ -489,13 +488,12 @@ Show status information about the Linspector instance
''') ''')
def do_about(self, text): def do_about(self, text):
self.print_color(GREEN, "Linspector Monitoring\n") self.print_color(GREEN, "Linspector System Monitoring\n")
self.print_color(YELLOW, "Developers:") self.print_color(YELLOW, "Developers:")
self.print_color(BLUE, " - Johannes Findeisen <hanez@linspector.org>") self.print_color(BLUE, " - Johannes Findeisen <hanez@linspector.org>")
self.print_color(BLUE, " - Rafael Timmerberg <ruff@linspector.org>\n") self.print_color(PURPLE, "(c) 2011 - 2015")
self.print_color(PURPLE, "(c) 2011 - 2013")
self.print_color(PURPLE, "Web: http://linspector.org") self.print_color(PURPLE, "Web: http://linspector.org")
self.print_color(PURPLE, "License: GNU Affero General Public License Version 3.0") self.print_color(PURPLE, "License: GNU General Public License Version 2")
def help_about(self): def help_about(self):
print(''' print('''

View file

@ -66,7 +66,7 @@ class DummyService(Service):
return False return False
def execute(self, execution): def execute(self, execution):
#print(str(self.sleep))
time.sleep(self.sleep) time.sleep(self.sleep)
d = {"Fail": str(self.fail), "Sleep": str(self.sleep)} d = {"Fail": str(self.fail), "Sleep": str(self.sleep)}