mirror of
https://github.com/django-q2/django-q2.git
synced 2026-09-19 19:08:07 +08:00
@@ -426,7 +426,6 @@ def save_cached(task, broker):
|
||||
iter_count = task.get('iter_count', 0)
|
||||
# if it's a group append to the group list
|
||||
if group:
|
||||
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 iter group, check if we are ready
|
||||
|
||||
@@ -16,35 +16,35 @@ from django_q.brokers import get_broker
|
||||
|
||||
def async(func, *args, **kwargs):
|
||||
"""Queue a task for the cluster."""
|
||||
# get options from q_options dict or direct from kwargs
|
||||
options = kwargs.pop('q_options', kwargs)
|
||||
broker = options.pop('broker', get_broker())
|
||||
# pop optionals
|
||||
opts = {'hook': None,
|
||||
'group': None,
|
||||
'save': None,
|
||||
'sync': None,
|
||||
'cached': Conf.CACHED,
|
||||
'iter_count': None,
|
||||
'iter_cached': None,
|
||||
'chain': None}
|
||||
for key in opts:
|
||||
opts[key] = options.pop(key, opts[key])
|
||||
# get an id
|
||||
keywords = kwargs.copy()
|
||||
opt_keys = ('hook', 'group', 'save', 'sync', 'cached', 'iter_count', 'iter_cached', 'chain', 'broker')
|
||||
q_options = keywords.pop('q_options', None)
|
||||
# get an id
|
||||
tag = uuid()
|
||||
# build the task package
|
||||
task = {'id': tag[1], 'name': tag[0],
|
||||
task = {'id': tag[1],
|
||||
'name': tag[0],
|
||||
'func': func,
|
||||
'args': args,
|
||||
'kwargs': kwargs,
|
||||
'started': timezone.now()}
|
||||
'args': args}
|
||||
# push optionals
|
||||
for key in opts:
|
||||
if opts[key] is not None:
|
||||
task[key] = opts[key]
|
||||
for key in opt_keys:
|
||||
if q_options and key in q_options:
|
||||
task[key] = q_options[key]
|
||||
elif key in keywords:
|
||||
task[key] = keywords.pop(key)
|
||||
# don't serialize the broker
|
||||
broker = task.pop('broker', get_broker())
|
||||
# overrides
|
||||
if 'cached' not in task and Conf.CACHED:
|
||||
task['cached'] = Conf.CACHED
|
||||
if 'sync' not in task and Conf.SYNC:
|
||||
task['sync'] = Conf.SYNC
|
||||
# finalize
|
||||
task['kwargs'] = keywords
|
||||
task['started'] = timezone.now()
|
||||
# sign it
|
||||
pack = signing.SignedPackage.dumps(task)
|
||||
if task.get('sync', False) or Conf.SYNC:
|
||||
if task.get('sync', False):
|
||||
return _sync(pack)
|
||||
# push it
|
||||
broker.enqueue(pack)
|
||||
@@ -106,7 +106,7 @@ def result(task_id, wait=0, cached=Conf.CACHED):
|
||||
r = Task.get_result(task_id)
|
||||
if r:
|
||||
return r
|
||||
if (time.time() - start) * 1000 >= wait:
|
||||
if (time.time() - start) * 1000 >= wait >= 0:
|
||||
break
|
||||
time.sleep(0.01)
|
||||
|
||||
@@ -122,7 +122,7 @@ def result_cached(task_id, wait=0, broker=None):
|
||||
r = broker.cache.get('{}:{}'.format(broker.list_key, task_id))
|
||||
if r:
|
||||
return signing.SignedPackage.loads(r)['result']
|
||||
if (time.time() - start) * 1000 >= wait:
|
||||
if (time.time() - start) * 1000 >= wait >= 0:
|
||||
break
|
||||
time.sleep(0.01)
|
||||
|
||||
@@ -142,14 +142,14 @@ def result_group(group_id, failures=False, wait=0, count=None, cached=Conf.CACHE
|
||||
start = time.time()
|
||||
if count:
|
||||
while True:
|
||||
if count_group(group_id) == count or wait and (time.time() - start) * 1000 >= wait:
|
||||
if count_group(group_id) == count or wait and (time.time() - start) * 1000 >= wait >= 0:
|
||||
break
|
||||
time.sleep(0.01)
|
||||
while True:
|
||||
r = Task.get_result_group(group_id, failures)
|
||||
if r:
|
||||
return r
|
||||
if (time.time() - start) * 1000 >= wait:
|
||||
if (time.time() - start) * 1000 >= wait >= 0:
|
||||
break
|
||||
time.sleep(0.01)
|
||||
|
||||
@@ -163,7 +163,7 @@ def result_group_cached(group_id, failures=False, wait=0, count=None, broker=Non
|
||||
start = time.time()
|
||||
if count:
|
||||
while True:
|
||||
if count_group_cached(group_id) == count or wait and (time.time() - start) * 1000 >= wait:
|
||||
if count_group_cached(group_id) == count or wait and (time.time() - start) * 1000 >= wait > 0:
|
||||
break
|
||||
time.sleep(0.01)
|
||||
while True:
|
||||
@@ -175,7 +175,7 @@ def result_group_cached(group_id, failures=False, wait=0, count=None, broker=Non
|
||||
if task['success'] or failures:
|
||||
result_list.append(task['result'])
|
||||
return result_list
|
||||
if (time.time() - start) * 1000 >= wait:
|
||||
if (time.time() - start) * 1000 >= wait >= 0:
|
||||
break
|
||||
time.sleep(0.01)
|
||||
|
||||
@@ -199,7 +199,7 @@ def fetch(task_id, wait=0, cached=Conf.CACHED):
|
||||
t = Task.get_task(task_id)
|
||||
if t:
|
||||
return t
|
||||
if (time.time() - start) * 1000 >= wait:
|
||||
if (time.time() - start) * 1000 >= wait >= 0:
|
||||
break
|
||||
time.sleep(0.01)
|
||||
|
||||
@@ -226,7 +226,7 @@ def fetch_cached(task_id, wait=0, broker=None):
|
||||
result=task['result'],
|
||||
success=task['success'])
|
||||
return t
|
||||
if (time.time() - start) * 1000 >= wait:
|
||||
if (time.time() - start) * 1000 >= wait >= 0:
|
||||
break
|
||||
time.sleep(0.01)
|
||||
|
||||
@@ -245,14 +245,14 @@ def fetch_group(group_id, failures=True, wait=0, count=None, cached=Conf.CACHED)
|
||||
start = time.time()
|
||||
if count:
|
||||
while True:
|
||||
if count_group(group_id) == count or wait and (time.time() - start) * 1000 >= wait:
|
||||
if count_group(group_id) == count or wait and (time.time() - start) * 1000 >= wait >= 0:
|
||||
break
|
||||
time.sleep(0.01)
|
||||
while True:
|
||||
r = Task.get_task_group(group_id, failures)
|
||||
if r:
|
||||
return r
|
||||
if (time.time() - start) * 1000 >= wait:
|
||||
if (time.time() - start) * 1000 >= wait >= 0:
|
||||
break
|
||||
time.sleep(0.01)
|
||||
|
||||
@@ -266,7 +266,7 @@ def fetch_group_cached(group_id, failures=True, wait=0, count=None, broker=None)
|
||||
start = time.time()
|
||||
if count:
|
||||
while True:
|
||||
if count_group_cached(group_id) == count or wait and (time.time() - start) * 1000 >= wait:
|
||||
if count_group_cached(group_id) == count or wait and (time.time() - start) * 1000 >= wait >= 0:
|
||||
break
|
||||
time.sleep(0.01)
|
||||
while True:
|
||||
@@ -289,7 +289,7 @@ def fetch_group_cached(group_id, failures=True, wait=0, count=None, broker=None)
|
||||
success=task['success'])
|
||||
task_list.append(t)
|
||||
return task_list
|
||||
if (time.time() - start) * 1000 >= wait:
|
||||
if (time.time() - start) * 1000 >= wait >= 0:
|
||||
break
|
||||
time.sleep(0.01)
|
||||
|
||||
@@ -472,7 +472,7 @@ class Iter(object):
|
||||
if self.started:
|
||||
return result(self.id, wait=wait, cached=self.cached)
|
||||
|
||||
def fetch(self, wait=0):
|
||||
def fetch(self, wait=0):
|
||||
"""
|
||||
get the task result objects.
|
||||
:param int wait: how many milliseconds to wait for a result
|
||||
@@ -560,6 +560,105 @@ class Chain(object):
|
||||
return len(self.chain)
|
||||
|
||||
|
||||
class Async(object):
|
||||
"""
|
||||
an async task
|
||||
"""
|
||||
|
||||
def __init__(self, func, *args, **kwargs):
|
||||
self.id = ''
|
||||
self.started = False
|
||||
self.func = func
|
||||
self.args = args
|
||||
self.kwargs = kwargs
|
||||
|
||||
@property
|
||||
def broker(self):
|
||||
return self._get_option('broker', None)
|
||||
|
||||
@broker.setter
|
||||
def broker(self, value):
|
||||
self._set_option('broker', value)
|
||||
|
||||
@property
|
||||
def sync(self):
|
||||
return self._get_option('sync', None)
|
||||
|
||||
@sync.setter
|
||||
def sync(self, value):
|
||||
self._set_option('sync', value)
|
||||
|
||||
@property
|
||||
def save(self):
|
||||
return self._get_option('save', None)
|
||||
|
||||
@save.setter
|
||||
def save(self, value):
|
||||
self._set_option('save', value)
|
||||
|
||||
@property
|
||||
def hook(self):
|
||||
return self._get_option('hook', None)
|
||||
|
||||
@hook.setter
|
||||
def hook(self, value):
|
||||
self._set_option('hook', value)
|
||||
|
||||
@property
|
||||
def group(self):
|
||||
return self._get_option('group', None)
|
||||
|
||||
@group.setter
|
||||
def group(self, value):
|
||||
self._set_option('group', value)
|
||||
|
||||
@property
|
||||
def cached(self):
|
||||
return self._get_option('cached', Conf.CACHED)
|
||||
|
||||
@cached.setter
|
||||
def cached(self, value):
|
||||
self._set_option('cached', value)
|
||||
|
||||
def _set_option(self, key, value):
|
||||
if 'q_options' in self.kwargs:
|
||||
self.kwargs['q_options'][key] = value
|
||||
else:
|
||||
self.kwargs[key] = value
|
||||
self.started = False
|
||||
|
||||
def _get_option(self, key, default=None):
|
||||
if 'q_options' in self.kwargs:
|
||||
return self.kwargs['q_options'].get(key, default)
|
||||
else:
|
||||
return self.kwargs.get(key, default)
|
||||
|
||||
def run(self):
|
||||
self.id = async(self.func, *self.args, **self.kwargs)
|
||||
self.started = True
|
||||
return self.id
|
||||
|
||||
def result(self, wait=0):
|
||||
|
||||
if self.started:
|
||||
return result(self.id, wait=wait, cached=self.cached)
|
||||
|
||||
def fetch(self, wait=0):
|
||||
|
||||
if self.started:
|
||||
return fetch(self.id, wait=wait, cached=self.cached)
|
||||
|
||||
def result_group(self, failures=False, wait=0, count=None):
|
||||
|
||||
if self.started and self.group:
|
||||
return result_group(self.group, failures=failures, wait=wait, count=count, cached=self.cached)
|
||||
|
||||
def fetch_group(self, failures=True, wait=0, count=None):
|
||||
|
||||
if self.started and self.group:
|
||||
return fetch_group(self.group, failures=failures, wait=wait, count=count, cached=self.cached)
|
||||
|
||||
|
||||
def _sync(pack):
|
||||
"""Simulate a package travelling through the cluster."""
|
||||
task_queue = Queue()
|
||||
|
||||
@@ -5,7 +5,7 @@ import pytest
|
||||
from django_q.cluster import pusher, worker, monitor
|
||||
from django_q.conf import Conf
|
||||
from django_q.tasks import async, result, fetch, count_group, result_group, fetch_group, delete_group, delete_cached, \
|
||||
async_iter, Chain, async_chain, Iter
|
||||
async_iter, Chain, async_chain, Iter, Async
|
||||
from django_q.brokers import get_broker
|
||||
|
||||
|
||||
@@ -145,3 +145,51 @@ def test_chain(broker):
|
||||
# test single
|
||||
rid = async_chain(['django_q.tests.tasks.hello', 'django_q.tests.tasks.hello'], sync=True, cached=True)
|
||||
assert result_group(rid, cached=True) == ['hello', 'hello']
|
||||
|
||||
|
||||
@pytest.mark.django_db
|
||||
def test_async_class(broker):
|
||||
broker.purge_queue()
|
||||
broker.cache.clear()
|
||||
a = Async('math.copysign')
|
||||
assert a.func == 'math.copysign'
|
||||
a.args = (1, -1)
|
||||
assert a.started is False
|
||||
a.cached = True
|
||||
assert a.cached is True
|
||||
a.sync = True
|
||||
assert a.sync is True
|
||||
a.broker = broker
|
||||
assert a.broker == broker
|
||||
a.run()
|
||||
assert a.result() == -1
|
||||
assert a.fetch().result == -1
|
||||
# again with kwargs
|
||||
a = Async('math.copysign', 1, -1, cached=True, sync=True, broker=broker)
|
||||
a.run()
|
||||
assert a.result() == -1
|
||||
# with q_options
|
||||
a = Async('math.copysign', 1, -1, q_options={'cached': True, 'sync': False, 'broker': broker})
|
||||
assert a.sync is False
|
||||
a.sync = True
|
||||
assert a.kwargs['q_options']['sync'] is True
|
||||
a.run()
|
||||
assert a.result() == -1
|
||||
a.group = 'async_class_test'
|
||||
assert a.group == 'async_class_test'
|
||||
a.save = False
|
||||
assert a.save is False
|
||||
a.hook = 'djq.tests.tasks.hello'
|
||||
assert a.hook == 'djq.tests.tasks.hello'
|
||||
assert a.started is False
|
||||
a.run()
|
||||
assert a.result_group() == [-1]
|
||||
assert a.fetch_group() == [a.fetch()]
|
||||
# global overrides
|
||||
Conf.SYNC = True
|
||||
Conf.CACHED = True
|
||||
a = Async('math.floor', 1.5)
|
||||
a.run()
|
||||
assert a.result() == 1
|
||||
Conf.SYNC = False
|
||||
Conf.CACHED = False
|
||||
|
||||
@@ -6,13 +6,12 @@
|
||||
#
|
||||
arrow==0.7.0
|
||||
blessed==1.12.0
|
||||
boto3==1.1.4
|
||||
botocore==1.2.11 # via boto3
|
||||
boto3==1.2.1
|
||||
botocore==1.3.1 # via boto3
|
||||
django-picklefield==0.3.2
|
||||
django-redis==4.2.0
|
||||
docutils==0.12 # via botocore
|
||||
future==0.15.2
|
||||
futures==2.2.0 # via boto3
|
||||
hiredis==0.2.0
|
||||
iron-core==1.1.9 # via iron-mq
|
||||
iron-mq==0.7
|
||||
|
||||
Reference in New Issue
Block a user