From f73ae66c0f2100c0057e216ec37f1d3fe9ce1721 Mon Sep 17 00:00:00 2001 From: Ilan Steemers Date: Sat, 3 Oct 2015 18:07:57 +0200 Subject: [PATCH] Adds async_iter command With async iter you can quickly run the same function on an iterable set of arguments. The results are held in the cache until all are done and collated into a database result. --- django_q/cluster.py | 39 ++++++++++++++++++++--------- django_q/tasks.py | 60 ++++++++++++++++++++++++++++++--------------- 2 files changed, 68 insertions(+), 31 deletions(-) diff --git a/django_q/cluster.py b/django_q/cluster.py index 8694294..4ca18f3 100644 --- a/django_q/cluster.py +++ b/django_q/cluster.py @@ -326,7 +326,7 @@ def monitor(result_queue, broker=None): broker.acknowledge(ack_id) # save the result if task.get('cached', False): - save_cache(task, broker) + save_cached(task, broker) else: save_task(task) # log the result @@ -413,22 +413,39 @@ def save_task(task): logger.error(e) -def save_cache(task, broker): - key = 'django_q:{}:results'.format(broker.list_key) +def save_cached(task, broker): + task_key = '{}:{}'.format(broker.list_key, task['id']) timeout = task['cached'] if timeout is True: timeout = None try: - task_package = signing.SignedPackage.dumps(task) group = task.get('group', False) + iter_count = task.get('iter_count', None) + # if it's a group append to the group list if group: - group_list = broker.cache.get('{}:{}'.format(key, group)) or [] - group_list.append(task_package) - broker.cache.set('{}:{}'.format(key, group), group_list, timeout) - else: - broker.cache.set('{}:{}'.format(key, task['id']), - task_package, - timeout) + task_key = '{}:{}:{}'.format(broker.list_key, group, task['id']) + group_key = '{}:{}:keys'.format(broker.list_key, group) + group_list = broker.cache.get(group_key) or [] + # if it's an inter group, check if we are ready + if iter_count and len(group_list) == iter_count-1: + group_args = '{}:{}:args'.format(broker.list_key, group) + # collate the results into a Task result + results = [signing.SignedPackage.loads(broker.cache.get(k))['result'] for k in group_list] + results.append(task['result']) + task['result'] = results + task['id'] = group + task['args'] = signing.SignedPackage.loads(broker.cache.get(group_args)) + save_task(task) + broker.cache.delete_many(group_list) + broker.cache.delete_many([group_key, group_args]) + return + # save the group list + group_list.append(task_key) + broker.cache.set(group_key, group_list) + # save the task + broker.cache.set(task_key, + signing.SignedPackage.dumps(task), + timeout) except Exception as e: logger.error(e) diff --git a/django_q/tasks.py b/django_q/tasks.py index 331d27c..17ae5b1 100644 --- a/django_q/tasks.py +++ b/django_q/tasks.py @@ -24,6 +24,7 @@ def async(func, *args, **kwargs): group = options.pop('group', None) save = options.pop('save', None) cached = options.pop('cached', Conf.CACHED) + iter_count = options.pop('iter_count', None) # get an id tag = uuid() # build the task package @@ -41,6 +42,8 @@ def async(func, *args, **kwargs): task['save'] = save if cached: task['cached'] = cached + if iter_count: + task['iter_count'] = iter_count # sign it pack = signing.SignedPackage.dumps(task) if sync or Conf.SYNC: @@ -116,10 +119,9 @@ def result_cached(task_id, wait=0, broker=None): """ if not broker: broker = get_broker() - key = 'django_q:{}:results'.format(broker.list_key) start = time.time() while True: - r = broker.cache.get('{}:{}'.format(key, task_id)) + r = broker.cache.get('{}:{}'.format(broker.list_key, task_id)) if r: return signing.SignedPackage.loads(r)['result'] if (time.time() - start) * 1000 >= wait: @@ -166,13 +168,12 @@ def result_group_cached(group_id, failures=False, wait=0, count=None, broker=Non if count_group_cached(group_id) == count or wait and (time.time() - start) * 1000 >= wait: break time.sleep(0.01) - key = 'django_q:{}:results'.format(broker.list_key) while True: - group_list = broker.cache.get('{}:{}'.format(key, group_id)) + group_list = broker.cache.get('{}:{}:keys'.format(broker.list_key, group_id)) if group_list: result_list = [] - for task_package in group_list: - task = signing.SignedPackage.loads(task_package) + for task_key in group_list: + task = signing.SignedPackage.loads(broker.cache.get(task_key)) if task['success'] or failures: result_list.append(task['result']) return result_list @@ -211,10 +212,9 @@ def fetch_cached(task_id, wait=0, broker=None): """ if not broker: broker = get_broker() - key = 'django_q:{}:results'.format(broker.list_key) start = time.time() while True: - r = broker.cache.get('{}:{}'.format(key, task_id)) + r = broker.cache.get('{}:{}'.format(broker.list_key, task_id)) if r: task = signing.SignedPackage.loads(r) t = Task(id=task['id'], @@ -271,13 +271,12 @@ def fetch_group_cached(group_id, failures=True, wait=0, count=None, broker=None) if count_group_cached(group_id) == count or wait and (time.time() - start) * 1000 >= wait: break time.sleep(0.01) - key = 'django_q:{}:results'.format(broker.list_key) while True: - group_list = broker.cache.get('{}:{}'.format(key, group_id)) + group_list = broker.cache.get('{}:{}:keys'.format(broker.list_key, group_id)) if group_list: task_list = [] - for task_package in group_list: - task = signing.SignedPackage.loads(task_package) + for task_key in group_list: + task = signing.SignedPackage.loads(broker.cache.get(task_key)) if task['success'] or failures: t = Task(id=task['id'], name=task['name'], @@ -318,14 +317,13 @@ def count_group_cached(group_id, failures=False, broker=None): """ if not broker: broker = get_broker() - key = 'django_q:{}:results'.format(broker.list_key) - group_list = broker.cache.get('{}:{}'.format(key, group_id)) + group_list = broker.cache.get('{}:{}:keys'.format(broker.list_key, group_id)) if group_list: if not failures: return len(group_list) failure_count = 0 - for task_package in group_list: - task = signing.SignedPackage.loads(task_package) + for task_key in group_list: + task = signing.SignedPackage.loads(broker.cache.get(task_key)) if not task['success']: failure_count += 1 return failure_count @@ -352,7 +350,10 @@ def delete_group_cached(group_id, broker=None): """ if not broker: broker = get_broker() - return delete_cached(group_id, broker) + group_key = '{}:{}:keys'.format(broker.list_key, group_id) + group_list = broker.cache.get(group_key) + broker.cache.delete_many(group_list) + broker.cache.delete(group_key) def delete_cached(task_id, broker=None): @@ -361,14 +362,13 @@ def delete_cached(task_id, broker=None): """ if not broker: broker = get_broker() - key = 'django_q:{}:results'.format(broker.list_key) - return broker.cache.delete('{}:{}'.format(key, task_id)) + return broker.cache.delete('{}:{}'.format(broker.list_key, task_id)) def queue_size(broker=None): """ Returns the current queue size. - Note that this doesn't count any tasks currently being processed by workers. + Note that this doesn't count any tasks curren key = 'django_q:{}:results'.format(broker.list_key)tly being processed by workers. :param broker: optional broker :return: current queue size @@ -379,6 +379,26 @@ def queue_size(broker=None): return broker.queue_size() +def async_iter(func, args_iter, **kwargs): + iter_count = len(args_iter) + iter_group = uuid()[1] + # clean up the kwargs + options = kwargs.get('q_options', kwargs) + options.pop('hook', None) + options['broker'] = options.get('broker', get_broker()) + options['group'] = iter_group + options['iter_count'] = iter_count + options['cached'] = True + # save the original arguments + broker = options['broker'] + broker.cache.set('{}:{}:args'.format(broker.list_key, iter_group), signing.SignedPackage.dumps(args_iter)) + for args in args_iter: + if type(args) is not tuple: + args = (args,) + async(func, *args, **options) + return iter_group + + def _sync(pack): """Simulate a package travelling through the cluster.""" task_queue = Queue()