mirror of
https://github.com/django-q2/django-q2.git
synced 2026-10-05 15:08:12 +08:00
Redis server now configurable with Q_REDIS. Moved constants to conf
This commit is contained in:
+28
-1
@@ -1,8 +1,17 @@
|
|||||||
from django.conf import settings
|
from signal import signal
|
||||||
from multiprocessing import cpu_count
|
from multiprocessing import cpu_count
|
||||||
|
|
||||||
|
from django.conf import settings
|
||||||
|
|
||||||
VERSION = '0.1.0'
|
VERSION = '0.1.0'
|
||||||
|
|
||||||
|
"""
|
||||||
|
Redis server connection
|
||||||
|
"""
|
||||||
|
try:
|
||||||
|
REDIS = settings.Q_REDIS
|
||||||
|
except AttributeError:
|
||||||
|
REDIS = {}
|
||||||
"""
|
"""
|
||||||
Prefixes the Redis keys. Defaults to django_q
|
Prefixes the Redis keys. Defaults to django_q
|
||||||
"""
|
"""
|
||||||
@@ -54,3 +63,21 @@ try:
|
|||||||
USE_TZ = settings.USE_TZ
|
USE_TZ = settings.USE_TZ
|
||||||
except AttributeError:
|
except AttributeError:
|
||||||
USE_TZ = False
|
USE_TZ = False
|
||||||
|
|
||||||
|
"""
|
||||||
|
Getting the SIGNAL names
|
||||||
|
"""
|
||||||
|
SIGNAL_NAMES = dict((getattr(signal, n), n) for n in dir(signal) if n.startswith('SIG') and '_' not in n)
|
||||||
|
|
||||||
|
"""
|
||||||
|
Redis List name
|
||||||
|
"""
|
||||||
|
Q_LIST = '{}:q'.format(PREFIX)
|
||||||
|
|
||||||
|
"""
|
||||||
|
Cluster status
|
||||||
|
"""
|
||||||
|
STARTING = 'Starting'
|
||||||
|
RUNNING = 'Running'
|
||||||
|
STOPPED = 'Stopped'
|
||||||
|
STOPPING = 'Stopping'
|
||||||
|
|||||||
+46
-12
@@ -4,7 +4,6 @@ from __future__ import print_function
|
|||||||
from __future__ import division
|
from __future__ import division
|
||||||
from __future__ import absolute_import
|
from __future__ import absolute_import
|
||||||
import ast
|
import ast
|
||||||
from builtins import dict
|
|
||||||
from builtins import range
|
from builtins import range
|
||||||
|
|
||||||
from future import standard_library
|
from future import standard_library
|
||||||
@@ -35,12 +34,11 @@ from django.core import signing
|
|||||||
from django.utils import timezone
|
from django.utils import timezone
|
||||||
|
|
||||||
# Local
|
# Local
|
||||||
from .conf import LOG_LEVEL, SECRET_KEY, SAVE_LIMIT, WORKERS, COMPRESSED, PREFIX
|
from .conf import LOG_LEVEL, SECRET_KEY, SAVE_LIMIT, WORKERS, COMPRESSED, PREFIX, REDIS, Q_LIST, SIGNAL_NAMES, STARTING, \
|
||||||
|
RUNNING, STOPPING, STOPPED
|
||||||
from .humanhash import uuid
|
from .humanhash import uuid
|
||||||
from .models import Task, Success, Schedule
|
from .models import Task, Success, Schedule
|
||||||
|
|
||||||
SIGNAL_NAMES = dict((getattr(signal, n), n) for n in dir(signal) if n.startswith('SIG') and '_' not in n)
|
|
||||||
|
|
||||||
logger = logging.getLogger('django-q')
|
logger = logging.getLogger('django-q')
|
||||||
|
|
||||||
|
|
||||||
@@ -55,20 +53,14 @@ except ImportError:
|
|||||||
# Set up standard logging handler
|
# Set up standard logging handler
|
||||||
if not logger.handlers:
|
if not logger.handlers:
|
||||||
logger.setLevel(level=getattr(logging, LOG_LEVEL))
|
logger.setLevel(level=getattr(logging, LOG_LEVEL))
|
||||||
|
|
||||||
formatter = logging.Formatter(fmt='%(asctime)s [django_q] %(message)s',
|
formatter = logging.Formatter(fmt='%(asctime)s [django_q] %(message)s',
|
||||||
datefmt='%H:%M:%S')
|
datefmt='%H:%M:%S')
|
||||||
handler = logging.StreamHandler()
|
handler = logging.StreamHandler()
|
||||||
handler.setFormatter(formatter)
|
handler.setFormatter(formatter)
|
||||||
logger.addHandler(handler)
|
logger.addHandler(handler)
|
||||||
|
|
||||||
Q_LIST = '{}:q'.format(PREFIX)
|
# Redis
|
||||||
STARTING = 'Starting'
|
r = redis.StrictRedis(**REDIS)
|
||||||
RUNNING = 'Running'
|
|
||||||
STOPPED = 'Stopped'
|
|
||||||
STOPPING = 'Stopping'
|
|
||||||
|
|
||||||
r = redis.StrictRedis()
|
|
||||||
|
|
||||||
|
|
||||||
class Cluster(object):
|
class Cluster(object):
|
||||||
@@ -264,6 +256,12 @@ class Sentinel(object):
|
|||||||
|
|
||||||
|
|
||||||
def pusher(task_queue, e, list_key=Q_LIST):
|
def pusher(task_queue, e, list_key=Q_LIST):
|
||||||
|
"""
|
||||||
|
Pulls tasks of the Redis List and puts them in the task queue
|
||||||
|
:type task_queue: multiprocessing.Queue
|
||||||
|
:type e: multiprocessing.Event
|
||||||
|
: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))
|
||||||
while True:
|
while True:
|
||||||
task = r.blpop(list_key, 1)
|
task = r.blpop(list_key, 1)
|
||||||
@@ -277,6 +275,10 @@ def pusher(task_queue, e, list_key=Q_LIST):
|
|||||||
|
|
||||||
|
|
||||||
def monitor(done_queue):
|
def monitor(done_queue):
|
||||||
|
"""
|
||||||
|
Gets finished tasks from the result queue and saves them to Django
|
||||||
|
:type done_queue: multiprocessing.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))
|
||||||
for task in iter(done_queue.get, 'STOP'):
|
for task in iter(done_queue.get, 'STOP'):
|
||||||
@@ -289,6 +291,11 @@ def monitor(done_queue):
|
|||||||
|
|
||||||
|
|
||||||
def worker(task_queue, done_queue):
|
def worker(task_queue, done_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 done_queue: multiprocessing.Queue
|
||||||
|
"""
|
||||||
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))
|
||||||
task = {}
|
task = {}
|
||||||
@@ -408,6 +415,10 @@ class PickleSerializer(object):
|
|||||||
|
|
||||||
|
|
||||||
class Status(object):
|
class Status(object):
|
||||||
|
"""
|
||||||
|
Cluster status base object
|
||||||
|
"""
|
||||||
|
|
||||||
def __init__(self, pid):
|
def __init__(self, pid):
|
||||||
self.workers = []
|
self.workers = []
|
||||||
self.tob = None
|
self.tob = None
|
||||||
@@ -424,6 +435,10 @@ class Status(object):
|
|||||||
|
|
||||||
|
|
||||||
class Stat(Status):
|
class Stat(Status):
|
||||||
|
"""
|
||||||
|
Status object for Cluster monitoring
|
||||||
|
"""
|
||||||
|
|
||||||
def __init__(self, sentinel, message=None):
|
def __init__(self, sentinel, message=None):
|
||||||
super(Stat, self).__init__(sentinel.parent_pid)
|
super(Stat, self).__init__(sentinel.parent_pid)
|
||||||
if message:
|
if message:
|
||||||
@@ -444,10 +459,17 @@ class Stat(Status):
|
|||||||
|
|
||||||
@property
|
@property
|
||||||
def key(self):
|
def key(self):
|
||||||
|
"""
|
||||||
|
:return: redis key for this cluster statistic
|
||||||
|
"""
|
||||||
return self.get_key(self.cluster_id)
|
return self.get_key(self.cluster_id)
|
||||||
|
|
||||||
@staticmethod
|
@staticmethod
|
||||||
def get_key(cluster_id):
|
def get_key(cluster_id):
|
||||||
|
"""
|
||||||
|
:param cluster_id: cluster ID
|
||||||
|
:return: redis key for the cluster statistic
|
||||||
|
"""
|
||||||
return '{}:cluster:{}'.format(PREFIX, cluster_id)
|
return '{}:cluster:{}'.format(PREFIX, cluster_id)
|
||||||
|
|
||||||
def save(self):
|
def save(self):
|
||||||
@@ -458,6 +480,11 @@ class Stat(Status):
|
|||||||
|
|
||||||
@staticmethod
|
@staticmethod
|
||||||
def get(cluster_id):
|
def get(cluster_id):
|
||||||
|
"""
|
||||||
|
gets the current status for the cluster
|
||||||
|
:param cluster_id: id of the cluster
|
||||||
|
:return: Stat or Status
|
||||||
|
"""
|
||||||
key = Stat.get_key(cluster_id)
|
key = Stat.get_key(cluster_id)
|
||||||
if r.exists(key):
|
if r.exists(key):
|
||||||
pack = r.get(key)
|
pack = r.get(key)
|
||||||
@@ -469,6 +496,10 @@ class Stat(Status):
|
|||||||
|
|
||||||
@staticmethod
|
@staticmethod
|
||||||
def get_all():
|
def get_all():
|
||||||
|
"""
|
||||||
|
Gets status for all currently running clusters with the same prefix and secret key
|
||||||
|
:return: Stat list
|
||||||
|
"""
|
||||||
stats = []
|
stats = []
|
||||||
keys = r.keys(pattern='{}:cluster:*'.format(PREFIX))
|
keys = r.keys(pattern='{}:cluster:*'.format(PREFIX))
|
||||||
if keys:
|
if keys:
|
||||||
@@ -482,6 +513,9 @@ class Stat(Status):
|
|||||||
|
|
||||||
|
|
||||||
def scheduler():
|
def scheduler():
|
||||||
|
"""
|
||||||
|
Creates a task from a schedule at the scheduled time and schedules next run
|
||||||
|
"""
|
||||||
for schedule in Schedule.objects.exclude(repeats=0).filter(next_run__lt=timezone.now()):
|
for schedule in Schedule.objects.exclude(repeats=0).filter(next_run__lt=timezone.now()):
|
||||||
args = ()
|
args = ()
|
||||||
kwargs = {}
|
kwargs = {}
|
||||||
|
|||||||
@@ -104,3 +104,12 @@ STATIC_URL = '/static/'
|
|||||||
|
|
||||||
# Django Q specific
|
# Django Q specific
|
||||||
Q_PREFIX = 'test_django_q'
|
Q_PREFIX = 'test_django_q'
|
||||||
|
|
||||||
|
Q_REDIS = {
|
||||||
|
#'host': 'localhost',
|
||||||
|
#'port': 6379,
|
||||||
|
# 'db': 1,
|
||||||
|
# 'socket_timeout': 3,
|
||||||
|
# 'password': '...',
|
||||||
|
# 'unix_socket_path': ''
|
||||||
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user