From 34e73a56f2492366f225fac67a030cb909168777 Mon Sep 17 00:00:00 2001 From: Johannes Findeisen Date: Wed, 11 Dec 2013 05:16:11 +0100 Subject: [PATCH] reverted to 7f9397eb30; needs some cleanups but then we are back online again.... ;) --- CHANGELOG | 1 - bin/linspector | 64 +++++++++++++++++++--------------------- examples/minimal.json | 5 ++-- linspector/core/job.py | 8 ++--- linspector/tasks/task.py | 36 +++++++++++----------- 5 files changed, 52 insertions(+), 62 deletions(-) diff --git a/CHANGELOG b/CHANGELOG index 710e666..109afee 100644 --- a/CHANGELOG +++ b/CHANGELOG @@ -1,7 +1,6 @@ 0.16 - Date ??? Next release is scheduled to christmas 2013 - Added a CHANGELOG file - - Added multi core support by instantiating multiple scheduler instances - Added "gui" command to lish providing an urwid (curses) based interface to Linspector - Added scheduling status to boot message diff --git a/bin/linspector b/bin/linspector index 9455ff1..d915f7a 100755 --- a/bin/linspector +++ b/bin/linspector @@ -20,7 +20,7 @@ along with this program. If not, see . """ -__version__ = "0.15.3/AMNESIA" +__version__ = "0.15/AMNESIA" import argparse import datetime @@ -29,10 +29,13 @@ import logging.handlers import os import os.path as path import sys +import multiprocessing +from multiprocessing import Pool from linspector.config.parser import FullConfigParser from linspector.core.interface import LinspectorInterface -from linspector.core.worker import LinspectorWorker +from linspector.core.job import LinspectorJob +from linspector.core.scheduler import LinspectorScheduler from linspector.frontends.lish import LishFrontend from linspector.tasks.task import TaskExecutor @@ -60,14 +63,11 @@ def parse_args(): parser.add_argument("-c", "--logcount", default=5, type=int, help="maximum number of logfiles in rotation (default: 5)") - parser.add_argument("-i", "--instances", default=1, type=int, - help="number of scheduler instances (default: 1)") - parser.add_argument("-m", "--logsize", default=10485760, type=int, help="maximum logfile size in bytes (default: 10485760)") - parser.add_argument("-t", "--threads", default=1000, type=int, - help="maximum number of scheduler threads (default: 1000)") + parser.add_argument("-t", "--threads", default=3500, type=int, + help="maximum number of scheduler threads (default: 3500)") parser.add_argument("-k", "--corethreads", default=0, type=int, help="number of scheduler core threads (default: 0)") @@ -91,6 +91,14 @@ def parse_args(): return parser.parse_args() +pool = Pool(processes=16) +print 'cpu_count() = %d\n' % multiprocessing.cpu_count() + + +def handle_job(job): + result = pool.apply_async(job.handle_call()) + + def main(): global lin_conf args = parse_args() @@ -114,26 +122,21 @@ def main(): logger.error(msg) exit() - executor = TaskExecutor() - #logger.debug("task executor instance :%s", str(task)) + scheduler = LinspectorScheduler({"apscheduler.threadpool.core_threads": args.corethreads, + "apscheduler.threadpool.max_threads": args.threads}) + scheduler.start() + TaskExecutor.Instance() job_count = 0 for layout in lin_conf.get_enabled_layouts(): for hostgroup in layout.get_hostgroups(): job_count += (hostgroup.get_services().__len__() * hostgroup.get_hosts().__len__()) - workers = [] - for i in range(0, args.instances): - process = LinspectorWorker(args.corethreads, args.threads) - workers.append(process) - process.daemon = True - start_date = datetime.datetime.now() time_delta = 0 jobs = [] count = 0 percent = 0 - instance = 0 for layout in lin_conf.get_enabled_layouts(): for hostgroup in layout.get_hostgroups(): for service in hostgroup.get_services(): @@ -151,22 +154,19 @@ def main(): time_delta += float(args.delay) new_start_date = start_date + datetime.timedelta(seconds=time_delta) - - worker = workers[instance] - worker.create_job(jobs, service, host, hostgroup, core, period, start_date=new_start_date, executor=executor) - - instance += 1 - if instance == args.instances: - instance = 0 + job = LinspectorJob(service, + host, + hostgroup.get_members(), + core, + hostgroup) + scheduler_job = period.createJob(scheduler, job, handle_job, start_date=new_start_date) + if scheduler_job is not None: + job.set_job(scheduler_job) + jobs.append(job) print("\nScheduled " + str(job_count) + " jobs") - schedulers = [] - for worker in workers: - schedulers.append(worker.get_scheduler()) - worker.start() - - interface = LinspectorInterface(jobs, schedulers, lin_conf, root_logger, __version__) + interface = LinspectorInterface(jobs, scheduler, lin_conf, root_logger, __version__) if "jsonrpc_backend" in core and core["jsonrpc_backend"]: from linspector.backends.jsonrpc import JsonrpcBackend @@ -184,15 +184,13 @@ def main(): if "shutdown_wait" in core: shutdown_wait = core["shutdown_wait"] - for worker in workers: - worker.shutdown(shutdown_wait) - worker.terminate() - + scheduler.shutdown(wait=shutdown_wait) logging.shutdown() if shutdown_wait: TaskExecutor.Instance().stop() else: TaskExecutor.Instance().stop_immediately() + if __name__ == "__main__": main() \ No newline at end of file diff --git a/examples/minimal.json b/examples/minimal.json index fc4dbb8..78defdc 100644 --- a/examples/minimal.json +++ b/examples/minimal.json @@ -21,8 +21,7 @@ "members": ["root"], "hosts": [ "a"], "services":[ - { "class": "etc/dummy", "args": { "sleep": 1, "fail": true }, "periods": ["fast"], "threshold": 1 }, - { "class": "etc/dummy", "args": { "sleep": 1, "fail": false }, "periods": ["fast"], "threshold": 1 } + { "class": "etc/dummy", "args": { "sleep": 1, "fail": true }, "periods": ["fast"], "threshold": 1 } ] } }, @@ -47,4 +46,4 @@ } } } -} +} \ No newline at end of file diff --git a/linspector/core/job.py b/linspector/core/job.py index 41fe9be..c0168b2 100644 --- a/linspector/core/job.py +++ b/linspector/core/job.py @@ -39,7 +39,7 @@ logger = getLogger(__name__) class LinspectorJob: - def __init__(self, service, host, members, core, hostgroup, executor): + def __init__(self, service, host, members, core, hostgroup): self.service = service self.host = host self.members = members @@ -49,7 +49,6 @@ class LinspectorJob: self.enabled = True self.scheduler_job = None self.job_id = self.hex_string() - self.executor = executor """ NONE job was not executed OK when everything is fine @@ -116,11 +115,8 @@ class LinspectorJob: for task in member.get_tasks(): if self.status.lower() in task.get_task_type().lower(): logger.debug("Executing Task of type: " + self.status) - try: + TaskExecutor.Instance().schedule_task(job_information, task) - self.executor.schedule_task(job_information, task) - except Exception, e: - logger.error("Error while executing: " + str(e)) def handle_call(self): logger.debug("handle call") logger.debug(self.service) diff --git a/linspector/tasks/task.py b/linspector/tasks/task.py index 281487f..d5d1e60 100644 --- a/linspector/tasks/task.py +++ b/linspector/tasks/task.py @@ -21,8 +21,6 @@ along with this program. If not, see . from logging import getLogger from threading import Event, Thread -from time import sleep -from Queue import Queue from linspector.utils.singleton import Singleton KEY_TYPE = "type" @@ -76,11 +74,11 @@ class Task(object): raise e +@Singleton class TaskExecutor(object): def __init__(self): - logger.debug("init taskExecutor") self.event = Event() - self.taskInfos = Queue() + self.taskInfos = [] task_thread = Thread(target=self._run_worker_thread) self._instantEnd = False self._running = True @@ -88,17 +86,20 @@ class TaskExecutor(object): task_thread.start() def _run_worker_thread(self): - logger.debug("start running taskExcecutor worker Thread") while self.is_running() or not self.is_instant_end(): + if len(self.taskInfos) == 0: + self.event.clear() + self.event.wait() + try: - logger.debug("try to get queued task from %s", self.taskInfos) - msg, task = self.taskInfos.get() - logger.debug("got task %s for msg: %s", str(task), str(msg)) - task.execute(msg) - self.taskInfos.task_done() + msg, task = self.taskInfos[0] + del self.taskInfos[0] + if task: + logger.debug("Starting Task Execution...") + task.execute(msg) + except Exception, e: - logger.error("Error: %s", e.message) - logger.debug("shutting down TaskExecutor!") + logger.error("Error " + str(e)) def is_instant_end(self): return self._instantEnd @@ -108,16 +109,13 @@ class TaskExecutor(object): def stop(self): self._running = False + self.event.set() def stop_immediately(self): self._running = False self._instantEnd = True + self.event.set() def schedule_task(self, msg, task): - try: - logger.debug("appending task '%s' for msg: %s", str(task), str(msg)) - self.taskInfos.put((msg, task)) - logger.debug("into queue: %s ", str(self.taskInfos)) - except Exception, e: - logger.debug("queue is probably full: %s", str(e)) - + self.taskInfos.append((msg, task)) + self.event.set()