mirror of
https://github.com/django-q2/django-q2.git
synced 2026-10-08 03:08:12 +08:00
Cleaned up some code
This commit is contained in:
+3
-11
@@ -246,7 +246,7 @@ class Sentinel(object):
|
|||||||
Stat(self, message).save()
|
Stat(self, message).save()
|
||||||
|
|
||||||
|
|
||||||
def pusher(task_queue, e, list_key=Q_LIST, r=None):
|
def pusher(task_queue, e, list_key=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
|
||||||
@@ -254,8 +254,6 @@ def pusher(task_queue, e, list_key=Q_LIST, r=None):
|
|||||||
: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))
|
||||||
if not r:
|
|
||||||
r = redis_client
|
|
||||||
while True:
|
while True:
|
||||||
task = r.blpop(list_key, 1)
|
task = r.blpop(list_key, 1)
|
||||||
if task:
|
if task:
|
||||||
@@ -291,9 +289,7 @@ def worker(task_queue, done_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_count = 0
|
task_count = 0
|
||||||
f = None
|
|
||||||
# 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
|
||||||
@@ -490,14 +486,12 @@ class Stat(Status):
|
|||||||
return self.done_q_size + self.task_q_size == 0
|
return self.done_q_size + self.task_q_size == 0
|
||||||
|
|
||||||
@staticmethod
|
@staticmethod
|
||||||
def get(cluster_id, r=None):
|
def get(cluster_id, r=redis_client):
|
||||||
"""
|
"""
|
||||||
gets the current status for the cluster
|
gets the current status for the cluster
|
||||||
:param cluster_id: id of the cluster
|
:param cluster_id: id of the cluster
|
||||||
:return: Stat or Status
|
:return: Stat or Status
|
||||||
"""
|
"""
|
||||||
if not r:
|
|
||||||
r = redis_client
|
|
||||||
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)
|
||||||
@@ -508,13 +502,11 @@ class Stat(Status):
|
|||||||
return Status(cluster_id)
|
return Status(cluster_id)
|
||||||
|
|
||||||
@staticmethod
|
@staticmethod
|
||||||
def get_all(r=None):
|
def get_all(r=redis_client):
|
||||||
"""
|
"""
|
||||||
Gets status for all currently running clusters with the same prefix and secret key
|
Gets status for all currently running clusters with the same prefix and secret key
|
||||||
:return: Stat list
|
:return: Stat list
|
||||||
"""
|
"""
|
||||||
if not r:
|
|
||||||
r = redis_client
|
|
||||||
stats = []
|
stats = []
|
||||||
keys = r.keys(pattern='{}:*'.format(Q_STAT))
|
keys = r.keys(pattern='{}:*'.format(Q_STAT))
|
||||||
if keys:
|
if keys:
|
||||||
|
|||||||
Reference in New Issue
Block a user