mirror of
https://github.com/django-q2/django-q2.git
synced 2026-10-04 10:08:12 +08:00
Adds Async class wrapper for the async function
This commit is contained in:
+112
-15
@@ -1,4 +1,5 @@
|
|||||||
"""Provides task functionality."""
|
"""Provides task functionality."""
|
||||||
|
from copy import deepcopy
|
||||||
from multiprocessing import Queue, Value
|
from multiprocessing import Queue, Value
|
||||||
|
|
||||||
# django
|
# django
|
||||||
@@ -17,31 +18,32 @@ from django_q.brokers import get_broker
|
|||||||
def async(func, *args, **kwargs):
|
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)
|
keywords = deepcopy(kwargs)
|
||||||
|
options = keywords.pop('q_options', keywords)
|
||||||
broker = options.pop('broker', get_broker())
|
broker = options.pop('broker', get_broker())
|
||||||
# pop optionals
|
# pop optionals
|
||||||
opts = {'hook': None,
|
opt_keys = {'hook': None,
|
||||||
'group': None,
|
'group': None,
|
||||||
'save': None,
|
'save': None,
|
||||||
'sync': None,
|
'sync': None,
|
||||||
'cached': Conf.CACHED,
|
'cached': Conf.CACHED,
|
||||||
'iter_count': None,
|
'iter_count': None,
|
||||||
'iter_cached': None,
|
'iter_cached': None,
|
||||||
'chain': None}
|
'chain': None}
|
||||||
for key in opts:
|
for key in opt_keys:
|
||||||
opts[key] = options.pop(key, opts[key])
|
opt_keys[key] = options.pop(key, opt_keys[key])
|
||||||
# get an id
|
# 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],
|
||||||
'func': func,
|
'func': func,
|
||||||
'args': args,
|
'args': args,
|
||||||
'kwargs': kwargs,
|
'kwargs': keywords,
|
||||||
'started': timezone.now()}
|
'started': timezone.now()}
|
||||||
# push optionals
|
# push optionals
|
||||||
for key in opts:
|
for key in opt_keys:
|
||||||
if opts[key] is not None:
|
if opt_keys[key] is not None:
|
||||||
task[key] = opts[key]
|
task[key] = opt_keys[key]
|
||||||
# sign it
|
# sign it
|
||||||
pack = signing.SignedPackage.dumps(task)
|
pack = signing.SignedPackage.dumps(task)
|
||||||
if task.get('sync', False) or Conf.SYNC:
|
if task.get('sync', False) or Conf.SYNC:
|
||||||
@@ -560,6 +562,101 @@ class Chain(object):
|
|||||||
return len(self.chain)
|
return len(self.chain)
|
||||||
|
|
||||||
|
|
||||||
|
class Async(object):
|
||||||
|
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):
|
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