mirror of
https://github.com/django-q2/django-q2.git
synced 2026-10-08 08:28:12 +08:00
Adds task chains
This commit is contained in:
+8
-2
@@ -393,6 +393,9 @@ def save_task(task):
|
|||||||
# SAVE LIMIT < 0 : Don't save success
|
# SAVE LIMIT < 0 : Don't save success
|
||||||
if not task.get('save', Conf.SAVE_LIMIT > 0) and task['success']:
|
if not task.get('save', Conf.SAVE_LIMIT > 0) and task['success']:
|
||||||
return
|
return
|
||||||
|
# async next in a chain
|
||||||
|
if task.get('chain', None):
|
||||||
|
tasks.async_chain(task['chain'], group=task['group'], cached=task['cached'])
|
||||||
# SAVE LIMIT > 0: Prune database, SAVE_LIMIT 0: No pruning
|
# SAVE LIMIT > 0: Prune database, SAVE_LIMIT 0: No pruning
|
||||||
db.close_old_connections()
|
db.close_old_connections()
|
||||||
try:
|
try:
|
||||||
@@ -419,14 +422,14 @@ def save_cached(task, broker):
|
|||||||
if timeout is True:
|
if timeout is True:
|
||||||
timeout = None
|
timeout = None
|
||||||
try:
|
try:
|
||||||
group = task.get('group', False)
|
group = task.get('group', None)
|
||||||
iter_count = task.get('iter_count', 0)
|
iter_count = task.get('iter_count', 0)
|
||||||
# if it's a group append to the group list
|
# if it's a group append to the group list
|
||||||
if group:
|
if group:
|
||||||
task_key = '{}:{}:{}'.format(broker.list_key, group, task['id'])
|
task_key = '{}:{}:{}'.format(broker.list_key, group, task['id'])
|
||||||
group_key = '{}:{}:keys'.format(broker.list_key, group)
|
group_key = '{}:{}:keys'.format(broker.list_key, group)
|
||||||
group_list = broker.cache.get(group_key) or []
|
group_list = broker.cache.get(group_key) or []
|
||||||
# if it's an inter group, check if we are ready
|
# if it's an iter group, check if we are ready
|
||||||
if iter_count and len(group_list) == iter_count-1:
|
if iter_count and len(group_list) == iter_count-1:
|
||||||
group_args = '{}:{}:args'.format(broker.list_key, group)
|
group_args = '{}:{}:args'.format(broker.list_key, group)
|
||||||
# collate the results into a Task result
|
# collate the results into a Task result
|
||||||
@@ -448,6 +451,9 @@ def save_cached(task, broker):
|
|||||||
# save the group list
|
# save the group list
|
||||||
group_list.append(task_key)
|
group_list.append(task_key)
|
||||||
broker.cache.set(group_key, group_list)
|
broker.cache.set(group_key, group_list)
|
||||||
|
# async next in a chain
|
||||||
|
if task.get('chain', None):
|
||||||
|
tasks.async_chain(task['chain'], group=group, cached=task['cached'])
|
||||||
# save the task
|
# save the task
|
||||||
broker.cache.set(task_key,
|
broker.cache.set(task_key,
|
||||||
signing.SignedPackage.dumps(task),
|
signing.SignedPackage.dumps(task),
|
||||||
|
|||||||
+1
-1
@@ -79,7 +79,7 @@ class Task(models.Model):
|
|||||||
return (self.stopped - self.started).total_seconds()
|
return (self.stopped - self.started).total_seconds()
|
||||||
|
|
||||||
def __unicode__(self):
|
def __unicode__(self):
|
||||||
return self.name
|
return u'{}'.format(self.name or self.id)
|
||||||
|
|
||||||
class Meta:
|
class Meta:
|
||||||
app_label = 'django_q'
|
app_label = 'django_q'
|
||||||
|
|||||||
+94
-21
@@ -18,15 +18,19 @@ def async(func, *args, **kwargs):
|
|||||||
"""Queue a task for the cluster."""
|
"""Queue a task for the cluster."""
|
||||||
# get options from q_options dict or direct from kwargs
|
# get options from q_options dict or direct from kwargs
|
||||||
options = kwargs.pop('q_options', kwargs)
|
options = kwargs.pop('q_options', kwargs)
|
||||||
hook = options.pop('hook', None)
|
|
||||||
broker = options.pop('broker', get_broker())
|
broker = options.pop('broker', get_broker())
|
||||||
sync = options.pop('sync', False)
|
sync = options.pop('sync', False)
|
||||||
group = options.pop('group', None)
|
# pop optionals
|
||||||
save = options.pop('save', None)
|
opts = {'hook': None,
|
||||||
cached = options.pop('cached', Conf.CACHED)
|
'group': None,
|
||||||
iter_count = options.pop('iter_count', None)
|
'save': None,
|
||||||
iter_cached = options.pop('iter_cached', None)
|
'cached': Conf.CACHED,
|
||||||
# get an id
|
'iter_count': None,
|
||||||
|
'iter_cached': None,
|
||||||
|
'chain': None}
|
||||||
|
for key in opts:
|
||||||
|
opts[key] = options.pop(key, opts[key])
|
||||||
|
# get an id
|
||||||
tag = uuid()
|
tag = uuid()
|
||||||
# build the task package
|
# build the task package
|
||||||
task = {'id': tag[1], 'name': tag[0],
|
task = {'id': tag[1], 'name': tag[0],
|
||||||
@@ -34,19 +38,10 @@ def async(func, *args, **kwargs):
|
|||||||
'args': args,
|
'args': args,
|
||||||
'kwargs': kwargs,
|
'kwargs': kwargs,
|
||||||
'started': timezone.now()}
|
'started': timezone.now()}
|
||||||
# add optionals
|
# push optionals
|
||||||
if hook:
|
for key in opts:
|
||||||
task['hook'] = hook
|
if opts[key] is not None:
|
||||||
if group:
|
task[key] = opts[key]
|
||||||
task['group'] = group
|
|
||||||
if save is not None:
|
|
||||||
task['save'] = save
|
|
||||||
if cached:
|
|
||||||
task['cached'] = cached
|
|
||||||
if iter_count:
|
|
||||||
task['iter_count'] = iter_count
|
|
||||||
if iter_cached:
|
|
||||||
task['iter_cached'] = iter_cached
|
|
||||||
# sign it
|
# sign it
|
||||||
pack = signing.SignedPackage.dumps(task)
|
pack = signing.SignedPackage.dumps(task)
|
||||||
if sync or Conf.SYNC:
|
if sync or Conf.SYNC:
|
||||||
@@ -371,7 +366,7 @@ def delete_cached(task_id, broker=None):
|
|||||||
def queue_size(broker=None):
|
def queue_size(broker=None):
|
||||||
"""
|
"""
|
||||||
Returns the current queue size.
|
Returns the current queue size.
|
||||||
Note that this doesn't count any tasks curren key = 'django_q:{}:results'.format(broker.list_key)tly being processed by workers.
|
Note that this doesn't count any tasks currently being processed by workers.
|
||||||
|
|
||||||
:param broker: optional broker
|
:param broker: optional broker
|
||||||
:return: current queue size
|
:return: current queue size
|
||||||
@@ -383,6 +378,9 @@ def queue_size(broker=None):
|
|||||||
|
|
||||||
|
|
||||||
def async_iter(func, args_iter, **kwargs):
|
def async_iter(func, args_iter, **kwargs):
|
||||||
|
"""
|
||||||
|
async a function with iterable arguments
|
||||||
|
"""
|
||||||
iter_count = len(args_iter)
|
iter_count = len(args_iter)
|
||||||
iter_group = uuid()[1]
|
iter_group = uuid()[1]
|
||||||
# clean up the kwargs
|
# clean up the kwargs
|
||||||
@@ -404,6 +402,81 @@ def async_iter(func, args_iter, **kwargs):
|
|||||||
return iter_group
|
return iter_group
|
||||||
|
|
||||||
|
|
||||||
|
def async_chain(chain, group=None, cached=Conf.CACHED, sync=Conf.SYNC):
|
||||||
|
"""
|
||||||
|
async a chain of tasks
|
||||||
|
the chain must be in the format [(func,(args),{kwargs}),(func,(args),{kwargs})]
|
||||||
|
"""
|
||||||
|
if not group:
|
||||||
|
group = uuid()[1]
|
||||||
|
args = ()
|
||||||
|
kwargs = {}
|
||||||
|
task = chain.pop(0)
|
||||||
|
if type(task) is not tuple:
|
||||||
|
task = (task,)
|
||||||
|
if len(task) > 1:
|
||||||
|
args = task[1]
|
||||||
|
if len(task) > 2:
|
||||||
|
kwargs = task[2]
|
||||||
|
kwargs['chain'] = chain
|
||||||
|
kwargs['group'] = group
|
||||||
|
kwargs['cached'] = cached
|
||||||
|
kwargs['sync'] = sync
|
||||||
|
async(task[0], *args, **kwargs)
|
||||||
|
return group
|
||||||
|
|
||||||
|
|
||||||
|
class Chain(object):
|
||||||
|
"""
|
||||||
|
A sequential chain of tasks
|
||||||
|
"""
|
||||||
|
def __init__(self, chain=None, group=None, cached=Conf.CACHED, sync=Conf.SYNC):
|
||||||
|
self.chain = chain or []
|
||||||
|
self.group = group or ''
|
||||||
|
self.cached = cached
|
||||||
|
self.sync = sync
|
||||||
|
|
||||||
|
def append(self, func, *args, **kwargs):
|
||||||
|
"""
|
||||||
|
add a task to the chain
|
||||||
|
takes the same parameters as async()
|
||||||
|
"""
|
||||||
|
task = (func, args, kwargs)
|
||||||
|
self.chain.append(task)
|
||||||
|
|
||||||
|
def run(self):
|
||||||
|
"""
|
||||||
|
Start queueing the chain to the worker cluster
|
||||||
|
:return: the chain's group id
|
||||||
|
"""
|
||||||
|
self.group = async_chain(self.chain, group=self.group, cached=self.cached, sync=self.sync)
|
||||||
|
return self.group
|
||||||
|
|
||||||
|
def result(self, wait=0):
|
||||||
|
"""
|
||||||
|
return the full list of results from the chain when it finishes. blocks until timeout.
|
||||||
|
:param int wait: how many milliseconds to wait for a result
|
||||||
|
:return: an unsorted list of results
|
||||||
|
"""
|
||||||
|
return result_group(self.group, wait=wait, count=len(self.chain), cached=self.cached)
|
||||||
|
|
||||||
|
def fetch(self, failures=True, wait=0):
|
||||||
|
"""
|
||||||
|
get the task result objects from the chain when it finishes. blocks until timeout.
|
||||||
|
:param failures: include failed tasks
|
||||||
|
:param int wait: how many milliseconds to wait for a result
|
||||||
|
:return: an unsorted list of task objects
|
||||||
|
"""
|
||||||
|
return fetch_group(self.group, failures=failures, wait=wait, count=len(self.chain), cached=self.cached)
|
||||||
|
|
||||||
|
def current(self):
|
||||||
|
"""
|
||||||
|
get the index of the currently executing chain element
|
||||||
|
:return int: current chain index
|
||||||
|
"""
|
||||||
|
return count_group(self.group, cached=self.cached)
|
||||||
|
|
||||||
|
|
||||||
def _sync(pack):
|
def _sync(pack):
|
||||||
"""Simulate a package travelling through the cluster."""
|
"""Simulate a package travelling through the cluster."""
|
||||||
task_queue = Queue()
|
task_queue = Queue()
|
||||||
|
|||||||
Reference in New Issue
Block a user