Adds timeout mechanism for workers

Workers now share a ctype with the sentinel to indicate if they are currently executing a task.  This timer is incremented each guard loop until the worker resets it after finishing a job or until it reaches the TIMEOUT value and the worker is terminated by the sentinel.
This commit is contained in:
Ilan Steemers
2015-07-01 10:54:36 +02:00
parent 5937f55fc5
commit 4e38925cea
4 changed files with 23 additions and 11 deletions
+3
View File
@@ -32,6 +32,9 @@ class Conf(object):
# Number of tasks each worker can handle before it gets recycled. Useful for releasing memory # Number of tasks each worker can handle before it gets recycled. Useful for releasing memory
RECYCLE = conf.get('recycle', 500) RECYCLE = conf.get('recycle', 500)
# Number of seconds to wait for a worker to finish.
TIMEOUT = conf.get('timeout', None)
# The Django Admin label for this app # The Django Admin label for this app
LABEL = conf.get('label', 'Django Q') LABEL = conf.get('label', 'Django Q')
+15 -6
View File
@@ -15,7 +15,7 @@ import importlib
import logging import logging
import os import os
import signal import signal
from multiprocessing import Queue, Event, Process, current_process from multiprocessing import Queue, Event, Process, Value, current_process
import socket import socket
import sys import sys
from time import sleep from time import sleep
@@ -181,6 +181,7 @@ class Sentinel(object):
p = Process(target=target, args=args) p = Process(target=target, args=args)
p.daemon = True p.daemon = True
if target == worker: if target == worker:
p.timer = args[2]
self.pool.append(p) self.pool.append(p)
p.start() p.start()
return p return p
@@ -189,7 +190,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.done_queue) self.spawn_process(worker, self.task_queue, self.done_queue, Value('b', -1))
def spawn_monitor(self): def spawn_monitor(self):
return self.spawn_process(monitor, self.done_queue) return self.spawn_process(monitor, self.done_queue)
@@ -203,7 +204,7 @@ class Sentinel(object):
logger.warn("reincarnated pusher after death of {}".format(pid)) logger.warn("reincarnated pusher after death of {}".format(pid))
else: else:
self.spawn_worker() self.spawn_worker()
logger.warn("reincarnated work worker after death of {}".format(pid)) logger.warn("reincarnated worker after death of {}".format(pid))
self.reincarnations += 1 self.reincarnations += 1
def spawn_cluster(self): def spawn_cluster(self):
@@ -223,11 +224,15 @@ class Sentinel(object):
# Guard loop. Runs at least once # Guard loop. Runs at least once
while not self.stop_event.is_set() or not counter: while not self.stop_event.is_set() or not counter:
# Check Workers # Check Workers
for p in list(self.pool): for p in self.pool:
if not p.is_alive(): # Are you alive?
if not p.is_alive() or (Conf.TIMEOUT and int(p.timer.value) >= Conf.TIMEOUT):
p.terminate() p.terminate()
self.pool.remove(p) self.pool.remove(p)
self.reincarnate(p.pid) self.reincarnate(p.pid)
# Increment timer if work is being done
if p.timer.value >= 0:
p.timer.value += 1
# Check Monitor # Check Monitor
if not self.monitor.is_alive(): if not self.monitor.is_alive():
self.reincarnate(self.monitor.pid) self.reincarnate(self.monitor.pid)
@@ -307,11 +312,12 @@ def monitor(done_queue):
logger.info("{} stopped monitoring results".format(name)) logger.info("{} stopped monitoring results".format(name))
def worker(task_queue, done_queue): def worker(task_queue, done_queue, timer):
""" """
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
:type done_queue: multiprocessing.Queue :type done_queue: multiprocessing.Queue
:type timer: multiprocessing.Value
""" """
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))
@@ -319,6 +325,7 @@ def worker(task_queue, done_queue):
# 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'):
result = None result = None
timer.value = -1 # Idle
task_count += 1 task_count += 1
# unpickle the task # unpickle the task
try: try:
@@ -340,6 +347,7 @@ def worker(task_queue, done_queue):
# We're still going # We're still going
if not result: if not result:
# execute the payload # execute the payload
timer.value = 0 # Busy
try: try:
res = f(*task['args'], **task['kwargs']) res = f(*task['args'], **task['kwargs'])
result = (res, True) result = (res, True)
@@ -350,6 +358,7 @@ def worker(task_queue, done_queue):
task['success'] = result[1] task['success'] = result[1]
task['stopped'] = timezone.now() task['stopped'] = timezone.now()
done_queue.put(task) done_queue.put(task)
timer.value = -1 # Idle
# Recycle # Recycle
if task_count == Conf.RECYCLE and task_queue.qsize() == 0: if task_count == Conf.RECYCLE and task_queue.qsize() == 0:
break break
+3 -3
View File
@@ -1,6 +1,6 @@
import sys import sys
import os import os
from multiprocessing import Queue, Event from multiprocessing import Queue, Event, Value
import pytest import pytest
@@ -74,7 +74,7 @@ def test_cluster(r):
assert r.llen(list_key) == 0 assert r.llen(list_key) == 0
# Test work # Test work
task_queue.put('STOP') task_queue.put('STOP')
worker(task_queue, result_queue) worker(task_queue, result_queue, Value('b', -1))
assert task_queue.qsize() == 0 assert task_queue.qsize() == 0
assert result_queue.qsize() == 1 assert result_queue.qsize() == 1
# Test monitor # Test monitor
@@ -131,7 +131,7 @@ def test_async(r):
task_queue.put('STOP') task_queue.put('STOP')
# let a worker handle them # let a worker handle them
result_queue = Queue() result_queue = Queue()
worker(task_queue, result_queue) worker(task_queue, result_queue, Value('b', -1))
assert result_queue.qsize() == task_count assert result_queue.qsize() == task_count
result_queue.put('STOP') result_queue.put('STOP')
# store the results # store the results
+2 -2
View File
@@ -1,4 +1,4 @@
from multiprocessing import Queue, Event from multiprocessing import Queue, Event, Value
import pytest import pytest
@@ -34,7 +34,7 @@ def test_scheduler(r):
task_queue.put('STOP') task_queue.put('STOP')
# let a worker handle them # let a worker handle them
result_queue = Queue() result_queue = Queue()
worker(task_queue, result_queue) worker(task_queue, result_queue, Value('b', -1))
assert result_queue.qsize() == 1 assert result_queue.qsize() == 1
result_queue.put('STOP') result_queue.put('STOP')
# store the results # store the results