closes db connection on worker and monitor spawn

This should prevent some problems with Postgresql where the workers would use stale db connections  and cause errors.
This commit is contained in:
Ilan
2015-07-27 12:31:38 +02:00
parent 82bd93d67f
commit ba58b4face
+64 -58
View File
@@ -31,6 +31,7 @@ except ImportError:
# Django # Django
from django.utils import timezone from django.utils import timezone
from django.utils.translation import ugettext_lazy as _ from django.utils.translation import ugettext_lazy as _
from django import db
# Local # Local
import signing import signing
@@ -328,6 +329,7 @@ def monitor(result_queue):
""" """
name = current_process().name name = current_process().name
logger.info(_("{} monitoring at {}").format(name, current_process().pid)) logger.info(_("{} monitoring at {}").format(name, current_process().pid))
db.close_old_connections()
for task in iter(result_queue.get, 'STOP'): for task in iter(result_queue.get, 'STOP'):
save_task(task) save_task(task)
if task['success']: if task['success']:
@@ -346,6 +348,7 @@ def worker(task_queue, result_queue, timer, timeout=Conf.TIMEOUT):
""" """
name = current_process().name name = current_process().name
logger.info(_('{} ready for work at {}').format(name, current_process().pid)) logger.info(_('{} ready for work at {}').format(name, current_process().pid))
db.close_old_connections()
task_count = 0 task_count = 0
# Start reading the task queue # Start reading the task queue
for pack in iter(task_queue.get, 'STOP'): for pack in iter(task_queue.get, 'STOP'):
@@ -399,9 +402,9 @@ def save_task(task):
if not task.get('save', Conf.SAVE_LIMIT > 0) and task['success']: if not task.get('save', Conf.SAVE_LIMIT > 0) and task['success']:
return return
# SAVE LIMIT > 0: Prune database, SAVE_LIMIT 0: No pruning # SAVE LIMIT > 0: Prune database, SAVE_LIMIT 0: No pruning
if task['success'] and 0 < Conf.SAVE_LIMIT < Success.objects.count():
Success.objects.last().delete()
try: try:
if task['success'] and 0 < Conf.SAVE_LIMIT < Success.objects.count():
Success.objects.last().delete()
Task.objects.create(id=task['id'], Task.objects.create(id=task['id'],
name=task['name'], name=task['name'],
func=task['func'], func=task['func'],
@@ -421,62 +424,65 @@ def scheduler(list_key=Conf.Q_LIST):
""" """
Creates a task from a schedule at the scheduled time and schedules next run Creates a task from a schedule at the scheduled time and schedules next run
""" """
for s in Schedule.objects.exclude(repeats=0).filter(next_run__lt=timezone.now()): try:
args = () for s in Schedule.objects.exclude(repeats=0).filter(next_run__lt=timezone.now()):
kwargs = {} args = ()
# get args, kwargs and hook kwargs = {}
if s.kwargs: # get args, kwargs and hook
try: if s.kwargs:
# eval should be safe here cause dict() try:
kwargs = eval('dict({})'.format(s.kwargs)) # eval should be safe here cause dict()
except SyntaxError: kwargs = eval('dict({})'.format(s.kwargs))
kwargs = {} except SyntaxError:
if s.args: kwargs = {}
args = ast.literal_eval(s.args) if s.args:
# single value won't eval to tuple, so: args = ast.literal_eval(s.args)
if type(args) != tuple: # single value won't eval to tuple, so:
args = (args,) if type(args) != tuple:
q_options = kwargs.get('q_options', {}) args = (args,)
if s.hook: q_options = kwargs.get('q_options', {})
q_options['hook'] = s.hook if s.hook:
# set up the next run time q_options['hook'] = s.hook
if not s.schedule_type == s.ONCE: # set up the next run time
next_run = arrow.get(s.next_run) if not s.schedule_type == s.ONCE:
if s.schedule_type == s.HOURLY: next_run = arrow.get(s.next_run)
next_run = next_run.replace(hours=+1) if s.schedule_type == s.HOURLY:
elif s.schedule_type == s.DAILY: next_run = next_run.replace(hours=+1)
next_run = next_run.replace(days=+1) elif s.schedule_type == s.DAILY:
elif s.schedule_type == s.WEEKLY: next_run = next_run.replace(days=+1)
next_run = next_run.replace(weeks=+1) elif s.schedule_type == s.WEEKLY:
elif s.schedule_type == s.MONTHLY: next_run = next_run.replace(weeks=+1)
next_run = next_run.replace(months=+1) elif s.schedule_type == s.MONTHLY:
elif s.schedule_type == s.QUARTERLY: next_run = next_run.replace(months=+1)
next_run = next_run.replace(months=+3) elif s.schedule_type == s.QUARTERLY:
elif s.schedule_type == s.YEARLY: next_run = next_run.replace(months=+3)
next_run = next_run.replace(years=+1) elif s.schedule_type == s.YEARLY:
s.next_run = next_run.datetime next_run = next_run.replace(years=+1)
s.repeats += -1 s.next_run = next_run.datetime
# send it to the cluster s.repeats += -1
q_options['list_key'] = list_key # send it to the cluster
q_options['group'] = s.name or s.id q_options['list_key'] = list_key
kwargs['q_options'] = q_options q_options['group'] = s.name or s.id
s.task = tasks.async(s.func, *args, **kwargs) kwargs['q_options'] = q_options
# log it s.task = tasks.async(s.func, *args, **kwargs)
if not s.task: # log it
logger.error( if not s.task:
_('{} failed to create a task from schedule [{}]').format(current_process().name, s.name or s.id)) logger.error(
else: _('{} failed to create a task from schedule [{}]').format(current_process().name, s.name or s.id))
logger.info( else:
_('{} created a task from schedule [{}]').format(current_process().name, s.name or s.id)) logger.info(
# default behavior is to delete a ONCE schedule _('{} created a task from schedule [{}]').format(current_process().name, s.name or s.id))
if s.schedule_type == s.ONCE: # default behavior is to delete a ONCE schedule
if s.repeats < 0: if s.schedule_type == s.ONCE:
s.delete() if s.repeats < 0:
return s.delete()
# but not if it has a positive repeats return
s.repeats = 0 # but not if it has a positive repeats
# save the schedule s.repeats = 0
s.save() # save the schedule
s.save()
except Exception as e:
logger.error(e)
def set_cpu_affinity(n, process_ids, actual=not Conf.TESTING): def set_cpu_affinity(n, process_ids, actual=not Conf.TESTING):