mirror of
https://github.com/django-q2/django-q2.git
synced 2026-10-07 02:08:13 +08:00
+9
-7
@@ -68,7 +68,8 @@ class Cluster(object):
|
|||||||
# Start Sentinel
|
# Start Sentinel
|
||||||
self.stop_event = Event()
|
self.stop_event = Event()
|
||||||
self.start_event = Event()
|
self.start_event = Event()
|
||||||
self.sentinel = Process(target=Sentinel, args=(self.stop_event, self.start_event, self.list_key, self.timeout))
|
self.sentinel = Process(target=Sentinel,
|
||||||
|
args=(self.stop_event, self.start_event, self.list_key, self.timeout))
|
||||||
self.sentinel.start()
|
self.sentinel.start()
|
||||||
logger.info(_('Q Cluster-{} starting.').format(self.pid))
|
logger.info(_('Q Cluster-{} starting.').format(self.pid))
|
||||||
while not self.start_event.is_set():
|
while not self.start_event.is_set():
|
||||||
@@ -87,7 +88,8 @@ class Cluster(object):
|
|||||||
return True
|
return True
|
||||||
|
|
||||||
def sig_handler(self, signum, frame):
|
def sig_handler(self, signum, frame):
|
||||||
logger.debug(_('{} got signal {}').format(current_process().name, Conf.SIGNAL_NAMES.get(signum, 'UNKNOWN')))
|
logger.debug(_('{} got signal {}').format(current_process().name,
|
||||||
|
Conf.SIGNAL_NAMES.get(signum, 'UNKNOWN')))
|
||||||
self.stop()
|
self.stop()
|
||||||
|
|
||||||
@property
|
@property
|
||||||
@@ -210,7 +212,7 @@ class Sentinel(object):
|
|||||||
self.pool = []
|
self.pool = []
|
||||||
Stat(self).save()
|
Stat(self).save()
|
||||||
# spawn worker pool
|
# spawn worker pool
|
||||||
for i in range(self.pool_size):
|
for _ in range(self.pool_size):
|
||||||
self.spawn_worker()
|
self.spawn_worker()
|
||||||
# spawn auxiliary
|
# spawn auxiliary
|
||||||
self.monitor = self.spawn_monitor()
|
self.monitor = self.spawn_monitor()
|
||||||
@@ -237,7 +239,7 @@ class Sentinel(object):
|
|||||||
continue
|
continue
|
||||||
# Decrement timer if work is being done
|
# Decrement timer if work is being done
|
||||||
if self.timeout and p.timer.value > 0:
|
if self.timeout and p.timer.value > 0:
|
||||||
p.timer.value -= cycle
|
p.timer.value -= cycle
|
||||||
# Check Monitor
|
# Check Monitor
|
||||||
if not self.monitor.is_alive():
|
if not self.monitor.is_alive():
|
||||||
self.reincarnate(self.monitor)
|
self.reincarnate(self.monitor)
|
||||||
@@ -295,11 +297,11 @@ class Sentinel(object):
|
|||||||
Stat(self).save()
|
Stat(self).save()
|
||||||
|
|
||||||
|
|
||||||
def pusher(task_queue, e, list_key=Conf.Q_LIST, r=redis_client):
|
def pusher(task_queue, event, list_key=Conf.Q_LIST, r=redis_client):
|
||||||
"""
|
"""
|
||||||
Pulls tasks of the Redis List and puts them in the task queue
|
Pulls tasks of the Redis List and puts them in the task queue
|
||||||
:type task_queue: multiprocessing.Queue
|
:type task_queue: multiprocessing.Queue
|
||||||
:type e: multiprocessing.Event
|
:type event: multiprocessing.Event
|
||||||
:type list_key: str
|
:type list_key: str
|
||||||
"""
|
"""
|
||||||
logger.info(_('{} pushing tasks at {}').format(current_process().name, current_process().pid))
|
logger.info(_('{} pushing tasks at {}').format(current_process().name, current_process().pid))
|
||||||
@@ -314,7 +316,7 @@ def pusher(task_queue, e, list_key=Conf.Q_LIST, r=redis_client):
|
|||||||
if task:
|
if task:
|
||||||
task_queue.put(task[1])
|
task_queue.put(task[1])
|
||||||
logger.debug(_('queueing from {}').format(list_key))
|
logger.debug(_('queueing from {}').format(list_key))
|
||||||
if e.is_set():
|
if event.is_set():
|
||||||
break
|
break
|
||||||
logger.info(_("{} stopped pushing tasks").format(current_process().name))
|
logger.info(_("{} stopped pushing tasks").format(current_process().name))
|
||||||
|
|
||||||
|
|||||||
@@ -8,6 +8,9 @@ import redis
|
|||||||
|
|
||||||
|
|
||||||
class Conf(object):
|
class Conf(object):
|
||||||
|
"""
|
||||||
|
Configuration class
|
||||||
|
"""
|
||||||
try:
|
try:
|
||||||
conf = settings.Q_CLUSTER
|
conf = settings.Q_CLUSTER
|
||||||
except AttributeError:
|
except AttributeError:
|
||||||
@@ -93,6 +96,10 @@ if Conf.DJANGO_REDIS:
|
|||||||
|
|
||||||
|
|
||||||
def get_redis_client():
|
def get_redis_client():
|
||||||
|
"""
|
||||||
|
Returns a connection from redis-py or django-redis
|
||||||
|
:return: a redis client
|
||||||
|
"""
|
||||||
if Conf.DJANGO_REDIS and django_redis:
|
if Conf.DJANGO_REDIS and django_redis:
|
||||||
return django_redis.get_redis_connection(Conf.DJANGO_REDIS)
|
return django_redis.get_redis_connection(Conf.DJANGO_REDIS)
|
||||||
return redis.StrictRedis(**Conf.REDIS)
|
return redis.StrictRedis(**Conf.REDIS)
|
||||||
|
|||||||
@@ -1,10 +1,5 @@
|
|||||||
from multiprocessing import Queue, Value
|
from multiprocessing import Queue, Value
|
||||||
|
|
||||||
try:
|
|
||||||
import cPickle as pickle
|
|
||||||
except ImportError:
|
|
||||||
import pickle
|
|
||||||
|
|
||||||
# django
|
# django
|
||||||
from django.utils import timezone
|
from django.utils import timezone
|
||||||
|
|
||||||
|
|||||||
Reference in New Issue
Block a user