mirror of
https://github.com/django-q2/django-q2.git
synced 2026-09-26 05:18:13 +08:00
Stabilize the stop procedure after forced termination of the monitor causes an unstable queue.
* added timeout to final stop * ony terminate on timeouts * renamed done_queue to result_queue
This commit is contained in:
+35
-22
@@ -129,10 +129,10 @@ class Sentinel(object):
|
||||
self.pool = []
|
||||
self.timeout = timeout
|
||||
self.task_queue = Queue()
|
||||
self.done_queue = Queue()
|
||||
self.result_queue = Queue()
|
||||
self.event_out = Event()
|
||||
self.monitor = None
|
||||
self.pusher = None
|
||||
self.monitor = Process()
|
||||
self.pusher = Process()
|
||||
if start:
|
||||
self.start()
|
||||
|
||||
@@ -144,7 +144,7 @@ class Sentinel(object):
|
||||
if not self.start_event.is_set() and not self.stop_event.is_set():
|
||||
return Conf.STARTING
|
||||
elif self.start_event.is_set() and not self.stop_event.is_set():
|
||||
if self.done_queue.qsize() == 0 and self.task_queue.qsize() == 0:
|
||||
if self.result_queue.qsize() == 0 and self.task_queue.qsize() == 0:
|
||||
return Conf.IDLE
|
||||
return Conf.WORKING
|
||||
elif self.stop_event.is_set() and self.start_event.is_set():
|
||||
@@ -174,17 +174,16 @@ class Sentinel(object):
|
||||
return self.spawn_process(pusher, self.task_queue, self.event_out, self.list_key, self.r)
|
||||
|
||||
def spawn_worker(self):
|
||||
self.spawn_process(worker, self.task_queue, self.done_queue, Value('b', -1))
|
||||
self.spawn_process(worker, self.task_queue, self.result_queue, Value('b', -1))
|
||||
|
||||
def spawn_monitor(self):
|
||||
return self.spawn_process(monitor, self.done_queue)
|
||||
return self.spawn_process(monitor, self.result_queue)
|
||||
|
||||
def reincarnate(self, process):
|
||||
"""
|
||||
:param process: the process to reincarnate
|
||||
:type process: Process or None
|
||||
:type process: Process
|
||||
"""
|
||||
process.terminate()
|
||||
if process == self.monitor:
|
||||
self.monitor = self.spawn_monitor()
|
||||
logger.error(_("reincarnated monitor {} after sudden death").format(process.name))
|
||||
@@ -195,11 +194,14 @@ class Sentinel(object):
|
||||
self.pool.remove(process)
|
||||
self.spawn_worker()
|
||||
if int(process.timer.value) >= self.timeout:
|
||||
# only need to terminate on timeout, otherwise we risk destabilizing the queues
|
||||
process.terminate()
|
||||
logger.warn(_("reincarnated worker {} after timeout").format(process.name))
|
||||
elif int(process.timer.value) == -2:
|
||||
logger.info(_("recycled worker {}").format(process.name))
|
||||
else:
|
||||
logger.error(_("reincarnated worker {} after death").format(process.name))
|
||||
|
||||
self.reincarnations += 1
|
||||
|
||||
def spawn_cluster(self):
|
||||
@@ -235,12 +237,12 @@ class Sentinel(object):
|
||||
self.reincarnate(self.pusher)
|
||||
# Call scheduler once a minute (or so)
|
||||
counter += 1
|
||||
if counter > 60:
|
||||
if counter > 120:
|
||||
counter = 0
|
||||
scheduler(list_key=self.list_key)
|
||||
# Save current status
|
||||
Stat(self).save()
|
||||
sleep(1)
|
||||
sleep(0.5)
|
||||
self.stop()
|
||||
|
||||
def stop(self):
|
||||
@@ -254,21 +256,32 @@ class Sentinel(object):
|
||||
sleep(0.2)
|
||||
Stat(self).save()
|
||||
# Put poison pills in the queue
|
||||
for _ in range(self.pool_size):
|
||||
for _ in range(len(self.pool)):
|
||||
self.task_queue.put('STOP')
|
||||
self.task_queue.close()
|
||||
# wait for the task queue to empty
|
||||
self.task_queue.join_thread()
|
||||
# Wait for all the workers to exit
|
||||
while len(self.pool) > 0:
|
||||
while len(self.pool):
|
||||
for p in self.pool:
|
||||
if not p.is_alive():
|
||||
logger.debug('{} stopped gracefully'.format(p.pid))
|
||||
self.pool.remove(p)
|
||||
sleep(0.2)
|
||||
Stat(self).save()
|
||||
# Finally stop the monitor
|
||||
self.done_queue.put('STOP')
|
||||
while self.status() != Conf.STOPPED:
|
||||
self.result_queue.put('STOP')
|
||||
self.result_queue.close()
|
||||
# Wait for the result queue to empty
|
||||
self.result_queue.join_thread()
|
||||
logger.info('{} waiting for the monitor.'.format(name))
|
||||
count = 0
|
||||
# Wait for everything to close or time out
|
||||
while self.status() == Conf.STOPPING and count < self.timeout * 5:
|
||||
sleep(0.2)
|
||||
Stat(self).save()
|
||||
count += 1
|
||||
# Final status
|
||||
Stat(self).save()
|
||||
self.pool = []
|
||||
|
||||
|
||||
@@ -291,27 +304,27 @@ def pusher(task_queue, e, list_key=Conf.Q_LIST, r=redis_client):
|
||||
logger.info(_("{} stopped pushing tasks").format(current_process().name))
|
||||
|
||||
|
||||
def monitor(done_queue):
|
||||
def monitor(result_queue):
|
||||
"""
|
||||
Gets finished tasks from the result queue and saves them to Django
|
||||
:type done_queue: multiprocessing.Queue
|
||||
:type result_queue: multiprocessing.Queue
|
||||
"""
|
||||
name = current_process().name
|
||||
logger.info(_("{} monitoring at {}").format(name, current_process().pid))
|
||||
for task in iter(done_queue.get, 'STOP'):
|
||||
for task in iter(result_queue.get, 'STOP'):
|
||||
save_task(task)
|
||||
if task['success']:
|
||||
logger.info(_("Processed [{}]").format(task['name']))
|
||||
else:
|
||||
logger.error(_("Failed [{}] - {}").format(task['name'], task['result']))
|
||||
save_task(task)
|
||||
logger.info(_("{} stopped monitoring results").format(name))
|
||||
|
||||
|
||||
def worker(task_queue, done_queue, timer):
|
||||
def worker(task_queue, result_queue, timer):
|
||||
"""
|
||||
Takes a task from the task queue, tries to execute it and puts the result back in the result queue
|
||||
:type task_queue: multiprocessing.Queue
|
||||
:type done_queue: multiprocessing.Queue
|
||||
:type result_queue: multiprocessing.Queue
|
||||
:type timer: multiprocessing.Value
|
||||
"""
|
||||
name = current_process().name
|
||||
@@ -352,7 +365,7 @@ def worker(task_queue, done_queue, timer):
|
||||
task['result'] = result[0]
|
||||
task['success'] = result[1]
|
||||
task['stopped'] = timezone.now()
|
||||
done_queue.put(task)
|
||||
result_queue.put(task)
|
||||
timer.value = -1 # Idle
|
||||
# Recycle
|
||||
if task_count == Conf.RECYCLE:
|
||||
|
||||
+4
-4
@@ -9,8 +9,8 @@ from django.utils import timezone
|
||||
from django.utils.translation import ugettext as _
|
||||
|
||||
# local
|
||||
from .conf import Conf, redis_client, logger
|
||||
from .tasks import SignedPackage
|
||||
from django_q.conf import Conf, redis_client, logger
|
||||
from django_q.tasks import SignedPackage
|
||||
|
||||
|
||||
def monitor(run_once=False):
|
||||
@@ -113,13 +113,13 @@ class Stat(Status):
|
||||
"""
|
||||
|
||||
def __init__(self, sentinel):
|
||||
super(Stat, self).__init__(sentinel.parent_pid)
|
||||
super(Stat, self).__init__(sentinel.parent_pid or sentinel.pid)
|
||||
self.r = sentinel.r
|
||||
self.tob = sentinel.tob
|
||||
self.reincarnations = sentinel.reincarnations
|
||||
self.sentinel = sentinel.pid
|
||||
self.status = sentinel.status()
|
||||
self.done_q_size = sentinel.done_queue.qsize()
|
||||
self.done_q_size = sentinel.result_queue.qsize()
|
||||
if sentinel.monitor:
|
||||
self.monitor = sentinel.monitor.pid
|
||||
self.task_q_size = sentinel.task_queue.qsize()
|
||||
|
||||
Reference in New Issue
Block a user