Merge pull request #18 from Koed00/dev

Adds a `timeout` override per task
This commit is contained in:
Ilan Steemers
2015-07-16 21:57:47 +02:00
4 changed files with 68 additions and 11 deletions
+9 -9
View File
@@ -176,7 +176,7 @@ class Sentinel(object):
return self.spawn_process(pusher, self.task_queue, self.event_out, self.list_key, self.r) return self.spawn_process(pusher, self.task_queue, self.event_out, self.list_key, self.r)
def spawn_worker(self): def spawn_worker(self):
self.spawn_process(worker, self.task_queue, self.result_queue, Value('b', -1)) self.spawn_process(worker, self.task_queue, self.result_queue, Value('b', -1), self.timeout)
def spawn_monitor(self): def spawn_monitor(self):
return self.spawn_process(monitor, self.result_queue) return self.spawn_process(monitor, self.result_queue)
@@ -195,7 +195,7 @@ class Sentinel(object):
else: else:
self.pool.remove(process) self.pool.remove(process)
self.spawn_worker() self.spawn_worker()
if self.timeout and int(process.timer.value) >= self.timeout: if self.timeout and int(process.timer.value) == 0:
# only need to terminate on timeout, otherwise we risk destabilizing the queues # only need to terminate on timeout, otherwise we risk destabilizing the queues
process.terminate() process.terminate()
logger.warn(_("reincarnated worker {} after timeout").format(process.name)) logger.warn(_("reincarnated worker {} after timeout").format(process.name))
@@ -231,12 +231,12 @@ class Sentinel(object):
# Check Workers # Check Workers
for p in self.pool: for p in self.pool:
# Are you alive? # Are you alive?
if not p.is_alive() or (self.timeout and int(p.timer.value) >= self.timeout): if not p.is_alive() or (self.timeout and int(p.timer.value) == 0):
self.reincarnate(p) self.reincarnate(p)
continue continue
# Increment timer if work is being done # Decrement timer if work is being done
if p.timer.value >= 0: if p.timer.value > 0:
p.timer.value += 1 p.timer.value -= 1
# Check Monitor # Check Monitor
if not self.monitor.is_alive(): if not self.monitor.is_alive():
self.reincarnate(self.monitor) self.reincarnate(self.monitor)
@@ -335,7 +335,7 @@ def monitor(result_queue):
logger.info(_("{} stopped monitoring results").format(name)) logger.info(_("{} stopped monitoring results").format(name))
def worker(task_queue, result_queue, timer): def worker(task_queue, result_queue, timer, timeout=Conf.TIMEOUT):
""" """
Takes a task from the task queue, tries to execute it and puts the result back in the result queue 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 task_queue: multiprocessing.Queue
@@ -370,7 +370,7 @@ def worker(task_queue, result_queue, timer):
# We're still going # We're still going
if not result: if not result:
# execute the payload # execute the payload
timer.value = 0 # Busy timer.value = task['kwargs'].pop('timeout', timeout or 0) # Busy
try: try:
res = f(*task['args'], **task['kwargs']) res = f(*task['args'], **task['kwargs'])
result = (res, True) result = (res, True)
@@ -384,7 +384,7 @@ def worker(task_queue, result_queue, timer):
timer.value = -1 # Idle timer.value = -1 # Idle
# Recycle # Recycle
if task_count == Conf.RECYCLE: if task_count == Conf.RECYCLE:
timer.value = -2 timer.value = -2 # Recycled
break break
logger.info(_('{} stopped doing work').format(name)) logger.info(_('{} stopped doing work').format(name))
+54
View File
@@ -213,6 +213,60 @@ def test_timeout(r):
assert start_event.is_set() assert start_event.is_set()
assert s.status() == Conf.STOPPED assert s.status() == Conf.STOPPED
assert s.reincarnations == 1 assert s.reincarnations == 1
r.delete(list_key)
@pytest.mark.django_db
def test_timeout(r):
# set up the Sentinel
list_key = 'timeout_test:q'
async('django_q.tests.tasks.count_forever', list_key=list_key)
start_event = Event()
stop_event = Event()
# Set a timer to stop the Sentinel
threading.Timer(3, stop_event.set).start()
s = Sentinel(stop_event, start_event, list_key=list_key, timeout=1)
assert start_event.is_set()
assert s.status() == Conf.STOPPED
assert s.reincarnations == 1
r.delete(list_key)
@pytest.mark.django_db
def test_timeout_override(r):
# set up the Sentinel
list_key = 'timeout_override_test:q'
async('django_q.tests.tasks.count_forever', list_key=list_key, timeout=1)
start_event = Event()
stop_event = Event()
# Set a timer to stop the Sentinel
threading.Timer(3, stop_event.set).start()
s = Sentinel(stop_event, start_event, list_key=list_key, timeout=10)
assert start_event.is_set()
assert s.status() == Conf.STOPPED
assert s.reincarnations == 1
r.delete(list_key)
@pytest.mark.django_db
def test_recycle(r):
# set up the Sentinel
list_key = 'test_recycle_test:q'
async('django_q.tests.tasks.multiply', 2, 2, list_key=list_key)
async('django_q.tests.tasks.multiply', 2, 2, list_key=list_key)
async('django_q.tests.tasks.multiply', 2, 2, list_key=list_key)
start_event = Event()
stop_event = Event()
# override settings
Conf.RECYCLE = 2
Conf.WORKERS = 1
# Set a timer to stop the Sentinel
threading.Timer(3, stop_event.set).start()
s = Sentinel(stop_event, start_event, list_key=list_key)
assert start_event.is_set()
assert s.status() == Conf.STOPPED
assert s.reincarnations == 1
r.delete(list_key)
@pytest.mark.django_db @pytest.mark.django_db
+3 -1
View File
@@ -65,11 +65,13 @@ recycle
The number of tasks a worker will process before recycling . Useful to release memory resources on a regular basis. Defaults to ``500``. The number of tasks a worker will process before recycling . Useful to release memory resources on a regular basis. Defaults to ``500``.
.. _timeout:
timeout timeout
~~~~~~~ ~~~~~~~
The number of seconds a worker is allowed to spend on a task before it's terminated. Defaults to ``None``, meaning it will never time out. The number of seconds a worker is allowed to spend on a task before it's terminated. Defaults to ``None``, meaning it will never time out.
Set this to something that makes sense for your project. Set this to something that makes sense for your project. Can be overridden for individual tasks.
compress compress
~~~~~~~~ ~~~~~~~~
+2 -1
View File
@@ -78,7 +78,7 @@ When you are making individual calls to :func:`async` a lot though, it can help
Reference Reference
--------- ---------
.. py:function:: async(func, *args, hook=None, sync=False, redis=None, **kwargs) .. py:function:: async(func, *args, hook=None, timeout=None, sync=False, redis=None, **kwargs)
Puts a task in the cluster queue Puts a task in the cluster queue
@@ -87,6 +87,7 @@ Reference
:type func: object :type func: object
:param hook: Optional function to call after execution :param hook: Optional function to call after execution
:type hook: object :type hook: object
:param int timeout: Overrides global cluster :ref:`timeout`.
:param bool sync: If set to True, async will simulate a task execution :param bool sync: If set to True, async will simulate a task execution
:param redis: Optional redis connection :param redis: Optional redis connection
:param kwargs: Keyword arguments for the task function :param kwargs: Keyword arguments for the task function