mirror of
https://github.com/django-q2/django-q2.git
synced 2026-10-04 14:08:11 +08:00
Merge pull request #36 from Koed00/dev
Moves unpacking task from Worker to Pusher
This commit is contained in:
+8
-8
@@ -315,7 +315,13 @@ def pusher(task_queue, event, list_key=Conf.Q_LIST, r=redis_client):
|
||||
sleep(10)
|
||||
break
|
||||
if task:
|
||||
task_queue.put(task[1])
|
||||
# unpack the task
|
||||
try:
|
||||
task = signing.SignedPackage.loads(task[1])
|
||||
except (TypeError, signing.BadSignature) as e:
|
||||
logger.error(e)
|
||||
continue
|
||||
task_queue.put(task)
|
||||
logger.debug(_('queueing from {}').format(list_key))
|
||||
if event.is_set():
|
||||
break
|
||||
@@ -351,16 +357,10 @@ def worker(task_queue, result_queue, timer, timeout=Conf.TIMEOUT):
|
||||
db.close_old_connections()
|
||||
task_count = 0
|
||||
# Start reading the task queue
|
||||
for pack in iter(task_queue.get, 'STOP'):
|
||||
for task in iter(task_queue.get, 'STOP'):
|
||||
result = None
|
||||
timer.value = -1 # Idle
|
||||
task_count += 1
|
||||
# unpickle the task
|
||||
try:
|
||||
task = signing.SignedPackage.loads(pack)
|
||||
except (TypeError, signing.BadSignature) as e:
|
||||
logger.error(e)
|
||||
continue
|
||||
# Get the function from the task
|
||||
logger.info(_('{} processing [{}]').format(name, task['name']))
|
||||
f = task['func']
|
||||
|
||||
+6
-5
@@ -40,7 +40,7 @@ def async(func, *args, **kwargs):
|
||||
# sign it
|
||||
pack = signing.SignedPackage.dumps(task)
|
||||
if sync:
|
||||
return _sync(task['id'], pack)
|
||||
return _sync(pack)
|
||||
# push it
|
||||
redis.rpush(list_key, pack)
|
||||
logger.debug('Pushed {}'.format(tag))
|
||||
@@ -150,13 +150,14 @@ def delete_group(group_id, tasks=False):
|
||||
return Task.delete_group(group_id, tasks)
|
||||
|
||||
|
||||
def _sync(task_id, pack):
|
||||
def _sync(pack):
|
||||
"""Simulate a package travelling through the cluster."""
|
||||
task_queue = Queue()
|
||||
result_queue = Queue()
|
||||
task_queue.put(pack)
|
||||
task = signing.SignedPackage.loads(pack)
|
||||
task_queue.put(task)
|
||||
task_queue.put('STOP')
|
||||
cluster.worker(task_queue, result_queue, Value('b', -1))
|
||||
cluster.worker(task_queue, result_queue, Value('f', -1))
|
||||
result_queue.put('STOP')
|
||||
cluster.monitor(result_queue)
|
||||
return task_id
|
||||
return task['id']
|
||||
|
||||
Reference in New Issue
Block a user