Add support for terminating local workers

This commit is contained in:
Roberto Rosario
2012-08-03 17:38:38 -04:00
parent ac60654cc9
commit 384209d2aa
7 changed files with 66 additions and 18 deletions

View File

@@ -13,6 +13,7 @@ from .literals import JOB_QUEUE_STATE_STARTED
LOCK_EXPIRE = 10
MAX_CPU_LOAD = 90.0
MAX_MEMORY_USAGE = 90.0
NODE_MAX_WORKERS = 1
logger = logging.getLogger(__name__)
@@ -20,27 +21,28 @@ logger = logging.getLogger(__name__)
def job_queue_poll():
logger.debug('starting')
node = Node.objects.myself()
if node.cpuload < MAX_CPU_LOAD and node.memory_usage < MAX_MEMORY_USAGE:
# Poll job queues if node is not overloaded
lock_id = u'job_queue_poll'
try:
lock = Lock.acquire_lock(lock_id, LOCK_EXPIRE)
except LockError:
pass
except Exception:
lock.release()
raise
else:
# Poll job queues if node is not overloaded
lock_id = u'job_queue_poll'
try:
lock = Lock.acquire_lock(lock_id, LOCK_EXPIRE)
except LockError:
pass
except Exception:
lock.release()
raise
else:
node = Node.objects.myself()
if node.cpuload < MAX_CPU_LOAD and node.memory_usage < MAX_MEMORY_USAGE and node.worker_set.count()<=NODE_MAX_WORKERS:
for job_queue in JobQueue.objects.filter(state=JOB_QUEUE_STATE_STARTED):
try:
job_item = job_queue.get_oldest_pending_job()
job_item.run()
except JobQueueNoPendingJobs:
logger.debug('no pending jobs for job queue: %s' % job_queue)
lock.release()
else:
logger.debug('CPU load or memory usage over limit')
else:
logger.debug('CPU load or memory usage over limit')
lock.release()
@simple_locking('house_keeping', 10)