Merge pull request #306 from P-EB/fix_304

Replaces async occurrences with alternatives
This commit is contained in:
Ilan Steemers
2018-08-01 14:11:21 +02:00
committed by GitHub
16 changed files with 190 additions and 190 deletions
+1 -1
View File
@@ -600,4 +600,4 @@
## [v0.1.0](https://github.com/koed00/django-q/tree/v0.1.0) (2015-06-28) ## [v0.1.0](https://github.com/koed00/django-q/tree/v0.1.0) (2015-06-28)
\* *This Change Log was automatically generated by [github_changelog_generator](https://github.com/skywinder/Github-Changelog-Generator)* \* *This Change Log was automatically generated by [github_changelog_generator](https://github.com/skywinder/Github-Changelog-Generator)*
+5 -5
View File
@@ -110,19 +110,19 @@ Check overall statistics with::
Creating Tasks Creating Tasks
~~~~~~~~~~~~~~ ~~~~~~~~~~~~~~
Use `async` from your code to quickly offload tasks: Use `enqueue` from your code to quickly offload tasks:
.. code:: python .. code:: python
from django_q.tasks import async, result from django_q.tasks import enqueue, result
# create the task # create the task
async('math.copysign', 2, -2) enqueue('math.copysign', 2, -2)
# or with a reference # or with a reference
import math.copysign import math.copysign
task_id = async(copysign, 2, -2) task_id = enqueue(copysign, 2, -2)
# get the result # get the result
task_result = result(task_id) task_result = result(task_id)
@@ -133,7 +133,7 @@ Use `async` from your code to quickly offload tasks:
# but in most cases you will want to use a hook: # but in most cases you will want to use a hook:
async('math.modf', 2.5, hook='hooks.print_result') enqueue('math.modf', 2.5, hook='hooks.print_result')
# hooks.py # hooks.py
def print_result(task): def print_result(task):
+1 -1
View File
@@ -12,7 +12,7 @@ default_app_config = 'django_q.apps.DjangoQConfig'
# root imports will slowly be deprecated. # root imports will slowly be deprecated.
# please import from the relevant sub modules # please import from the relevant sub modules
if django.VERSION[:2] < (1, 9): if django.VERSION[:2] < (1, 9):
from .tasks import async, schedule, result, result_group, fetch, fetch_group, count_group, delete_group, queue_size from .tasks import enqueue, schedule, result, result_group, fetch, fetch_group, count_group, delete_group, queue_size
from .models import Task, Schedule, Success, Failure from .models import Task, Schedule, Success, Failure
from .cluster import Cluster from .cluster import Cluster
from .status import Stat from .status import Stat
+2 -2
View File
@@ -2,7 +2,7 @@
from django.contrib import admin from django.contrib import admin
from django.utils.translation import ugettext_lazy as _ from django.utils.translation import ugettext_lazy as _
from django_q.tasks import async from django_q.tasks import enqueue
from django_q.models import Success, Failure, Schedule, OrmQ from django_q.models import Success, Failure, Schedule, OrmQ
from django_q.conf import Conf from django_q.conf import Conf
@@ -41,7 +41,7 @@ class TaskAdmin(admin.ModelAdmin):
def retry_failed(FailAdmin, request, queryset): def retry_failed(FailAdmin, request, queryset):
"""Submit selected tasks back to the queue.""" """Submit selected tasks back to the queue."""
for task in queryset: for task in queryset:
async(task.func, *task.args or (), hook=task.hook, **task.kwargs or {}) enqueue(task.func, *task.args or (), hook=task.hook, **task.kwargs or {})
task.delete() task.delete()
+5 -5
View File
@@ -405,9 +405,9 @@ def save_task(task, broker):
# 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 # enqueues next in a chain
if task.get('chain', None): if task.get('chain', None):
tasks.async_chain(task['chain'], group=task['group'], cached=task['cached'], sync=task['sync'], broker=broker) tasks.enqueue_chain(task['chain'], group=task['group'], cached=task['cached'], sync=task['sync'], broker=broker)
# 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:
@@ -473,9 +473,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, timeout) broker.cache.set(group_key, group_list, timeout)
# async next in a chain # enqueue next in a chain
if task.get('chain', None): if task.get('chain', None):
tasks.async_chain(task['chain'], group=group, cached=task['cached'], sync=task['sync'], broker=broker) tasks.enqueue_chain(task['chain'], group=group, cached=task['cached'], sync=task['sync'], broker=broker)
# save the task # save the task
broker.cache.set(task_key, broker.cache.set(task_key,
SignedPackage.dumps(task), SignedPackage.dumps(task),
@@ -536,7 +536,7 @@ def scheduler(broker=None):
q_options['broker'] = broker q_options['broker'] = broker
q_options['group'] = q_options.get('group', s.name or s.id) q_options['group'] = q_options.get('group', s.name or s.id)
kwargs['q_options'] = q_options kwargs['q_options'] = q_options
s.task = tasks.async(s.func, *args, **kwargs) s.task = tasks.enqueue(s.func, *args, **kwargs)
# log it # log it
if not s.task: if not s.task:
logger.error( logger.error(
+13 -13
View File
@@ -17,7 +17,7 @@ from django_q.signals import pre_enqueue
from django_q.queues import Queue from django_q.queues import Queue
def async(func, *args, **kwargs): def enqueue(func, *args, **kwargs):
"""Queue a task for the cluster.""" """Queue a task for the cluster."""
keywords = kwargs.copy() keywords = kwargs.copy()
opt_keys = ('hook', 'group', 'save', 'sync', 'cached', 'ack_failure', 'iter_count', 'iter_cached', 'chain', 'broker') opt_keys = ('hook', 'group', 'save', 'sync', 'cached', 'ack_failure', 'iter_count', 'iter_cached', 'chain', 'broker')
@@ -390,9 +390,9 @@ def queue_size(broker=None):
return broker.queue_size() return broker.queue_size()
def async_iter(func, args_iter, **kwargs): def enqueue_iter(func, args_iter, **kwargs):
""" """
async a function with iterable arguments enqueues a function with iterable arguments
""" """
iter_count = len(args_iter) iter_count = len(args_iter)
iter_group = uuid()[1] iter_group = uuid()[1]
@@ -411,13 +411,13 @@ def async_iter(func, args_iter, **kwargs):
for args in args_iter: for args in args_iter:
if type(args) is not tuple: if type(args) is not tuple:
args = (args,) args = (args,)
async(func, *args, **options) enqueue(func, *args, **options)
return iter_group return iter_group
def async_chain(chain, group=None, cached=Conf.CACHED, sync=Conf.SYNC, broker=None): def enqueue_chain(chain, group=None, cached=Conf.CACHED, sync=Conf.SYNC, broker=None):
""" """
async a chain of tasks enqueues a chain of tasks
the chain must be in the format [(func,(args),{kwargs}),(func,(args),{kwargs})] the chain must be in the format [(func,(args),{kwargs}),(func,(args),{kwargs})]
""" """
if not group: if not group:
@@ -436,7 +436,7 @@ def async_chain(chain, group=None, cached=Conf.CACHED, sync=Conf.SYNC, broker=No
kwargs['cached'] = cached kwargs['cached'] = cached
kwargs['sync'] = sync kwargs['sync'] = sync
kwargs['broker'] = broker or get_broker() kwargs['broker'] = broker or get_broker()
async(task[0], *args, **kwargs) enqueue(task[0], *args, **kwargs)
return group return group
@@ -472,7 +472,7 @@ class Iter(object):
self.kwargs['cached'] = self.cached self.kwargs['cached'] = self.cached
self.kwargs['sync'] = self.sync self.kwargs['sync'] = self.sync
self.kwargs['broker'] = self.broker self.kwargs['broker'] = self.broker
self.id = async_iter(self.func, self.args, **self.kwargs) self.id = enqueue_iter(self.func, self.args, **self.kwargs)
self.started = True self.started = True
return self.id return self.id
@@ -518,7 +518,7 @@ class Chain(object):
def append(self, func, *args, **kwargs): def append(self, func, *args, **kwargs):
""" """
add a task to the chain add a task to the chain
takes the same parameters as async() takes the same parameters as enqueue()
""" """
self.chain.append((func, args, kwargs)) self.chain.append((func, args, kwargs))
# remove existing results # remove existing results
@@ -532,8 +532,8 @@ class Chain(object):
Start queueing the chain to the worker cluster Start queueing the chain to the worker cluster
:return: the chain's group id :return: the chain's group id
""" """
self.group = async_chain(chain=self.chain[:], group=self.group, cached=self.cached, sync=self.sync, self.group = enqueue_chain(chain=self.chain[:], group=self.group, cached=self.cached, sync=self.sync,
broker=self.broker) broker=self.broker)
self.started = True self.started = True
return self.group return self.group
@@ -573,7 +573,7 @@ class Chain(object):
return len(self.chain) return len(self.chain)
class Async(object): class AsyncTask(object):
""" """
an async task an async task
""" """
@@ -647,7 +647,7 @@ class Async(object):
return self.kwargs.get(key, default) return self.kwargs.get(key, default)
def run(self): def run(self):
self.id = async(self.func, *self.args, **self.kwargs) self.id = enqueue(self.func, *self.args, **self.kwargs)
self.started = True self.started = True
return self.id return self.id
+20 -20
View File
@@ -5,8 +5,8 @@ import pytest
from django_q.cluster import pusher, worker, monitor from django_q.cluster import pusher, worker, monitor
from django_q.compat import range from django_q.compat import range
from django_q.conf import Conf from django_q.conf import Conf
from django_q.tasks import async, result, fetch, count_group, result_group, fetch_group, delete_group, delete_cached, \ from django_q.tasks import enqueue, result, fetch, count_group, result_group, fetch_group, delete_group, delete_cached, \
async_iter, Chain, async_chain, Iter, Async enqueue_iter, Chain, enqueue_chain, Iter, AsyncTask
from django_q.brokers import get_broker from django_q.brokers import get_broker
from django_q.queues import Queue from django_q.queues import Queue
@@ -23,14 +23,14 @@ def test_cached(broker):
broker.cache.clear() broker.cache.clear()
group = 'cache_test' group = 'cache_test'
# queue the tests # queue the tests
task_id = async('math.copysign', 1, -1, cached=True, broker=broker) task_id = enqueue('math.copysign', 1, -1, cached=True, broker=broker)
async('math.copysign', 1, -1, cached=True, broker=broker, group=group) enqueue('math.copysign', 1, -1, cached=True, broker=broker, group=group)
async('math.copysign', 1, -1, cached=True, broker=broker, group=group) enqueue('math.copysign', 1, -1, cached=True, broker=broker, group=group)
async('math.copysign', 1, -1, cached=True, broker=broker, group=group) enqueue('math.copysign', 1, -1, cached=True, broker=broker, group=group)
async('math.copysign', 1, -1, cached=True, broker=broker, group=group) enqueue('math.copysign', 1, -1, cached=True, broker=broker, group=group)
async('math.copysign', 1, -1, cached=True, broker=broker, group=group) enqueue('math.copysign', 1, -1, cached=True, broker=broker, group=group)
async('math.popysign', 1, -1, cached=True, broker=broker, group=group) enqueue('math.popysign', 1, -1, cached=True, broker=broker, group=group)
iter_id = async_iter('math.floor', [i for i in range(10)], cached=True) iter_id = enqueue_iter('math.floor', [i for i in range(10)], cached=True)
# test wait on cache # test wait on cache
# test wait timeout # test wait timeout
assert result(task_id, wait=10, cached=True) is None assert result(task_id, wait=10, cached=True) is None
@@ -86,10 +86,10 @@ def test_iter(broker):
it = [i for i in range(10)] it = [i for i in range(10)]
it2 = [(1, -1), (2, -1), (3, -4), (5, 6)] it2 = [(1, -1), (2, -1), (3, -4), (5, 6)]
it3 = (1, 2, 3, 4, 5) it3 = (1, 2, 3, 4, 5)
t = async_iter('math.floor', it, sync=True) t = enqueue_iter('math.floor', it, sync=True)
t2 = async_iter('math.copysign', it2, sync=True) t2 = enqueue_iter('math.copysign', it2, sync=True)
t3 = async_iter('math.floor', it3, sync=True) t3 = enqueue_iter('math.floor', it3, sync=True)
t4 = async_iter('math.floor', (1,), sync=True) t4 = enqueue_iter('math.floor', (1,), sync=True)
result_t = result(t) result_t = result(t)
assert result_t is not None assert result_t is not None
task_t = fetch(t) task_t = fetch(t)
@@ -140,15 +140,15 @@ def test_chain(broker):
t = task_chain.fetch() t = task_chain.fetch()
assert len(t) == task_chain.length() assert len(t) == task_chain.length()
# test single # test single
rid = async_chain(['django_q.tests.tasks.hello', 'django_q.tests.tasks.hello'], sync=True, cached=True) rid = enqueue_chain(['django_q.tests.tasks.hello', 'django_q.tests.tasks.hello'], sync=True, cached=True)
assert result_group(rid, cached=True) == ['hello', 'hello'] assert result_group(rid, cached=True) == ['hello', 'hello']
@pytest.mark.django_db @pytest.mark.django_db
def test_async_class(broker, monkeypatch): def test_asynctask_class(broker, monkeypatch):
broker.purge_queue() broker.purge_queue()
broker.cache.clear() broker.cache.clear()
a = Async('math.copysign') a = AsyncTask('math.copysign')
assert a.func == 'math.copysign' assert a.func == 'math.copysign'
a.args = (1, -1) a.args = (1, -1)
assert a.started is False assert a.started is False
@@ -162,11 +162,11 @@ def test_async_class(broker, monkeypatch):
assert a.result() == -1 assert a.result() == -1
assert a.fetch().result == -1 assert a.fetch().result == -1
# again with kwargs # again with kwargs
a = Async('math.copysign', 1, -1, cached=True, sync=True, broker=broker) a = AsyncTask('math.copysign', 1, -1, cached=True, sync=True, broker=broker)
a.run() a.run()
assert a.result() == -1 assert a.result() == -1
# with q_options # with q_options
a = Async('math.copysign', 1, -1, q_options={'cached': True, 'sync': False, 'broker': broker}) a = AsyncTask('math.copysign', 1, -1, q_options={'cached': True, 'sync': False, 'broker': broker})
assert a.sync is False assert a.sync is False
a.sync = True a.sync = True
assert a.kwargs['q_options']['sync'] is True assert a.kwargs['q_options']['sync'] is True
@@ -185,6 +185,6 @@ def test_async_class(broker, monkeypatch):
# global overrides # global overrides
monkeypatch.setattr(Conf, 'SYNC', True) monkeypatch.setattr(Conf, 'SYNC', True)
monkeypatch.setattr(Conf, 'CACHED', True) monkeypatch.setattr(Conf, 'CACHED', True)
a = Async('math.floor', 1.5) a = AsyncTask('math.floor', 1.5)
a.run() a.run()
assert a.result() == 1 assert a.result() == 1
+27 -27
View File
@@ -13,7 +13,7 @@ sys.path.insert(0, myPath + '/../')
from django_q.cluster import Cluster, Sentinel, pusher, worker, monitor, save_task from django_q.cluster import Cluster, Sentinel, pusher, worker, monitor, save_task
from django_q.compat import range from django_q.compat import range
from django_q.humanhash import DEFAULT_WORDLIST, uuid from django_q.humanhash import DEFAULT_WORDLIST, uuid
from django_q.tasks import fetch, fetch_group, async, result, result_group, count_group, delete_group, queue_size from django_q.tasks import fetch, fetch_group, enqueue, result, result_group, count_group, delete_group, queue_size
from django_q.models import Task, Success from django_q.models import Task, Success
from django_q.conf import Conf from django_q.conf import Conf
from django_q.status import Stat from django_q.status import Stat
@@ -42,7 +42,7 @@ def test_redis_connection(broker):
@pytest.mark.django_db @pytest.mark.django_db
def test_sync(broker): def test_sync(broker):
task = async('django_q.tests.tasks.count_letters', DEFAULT_WORDLIST, broker=broker, sync=True) task = enqueue('django_q.tests.tasks.count_letters', DEFAULT_WORDLIST, broker=broker, sync=True)
assert result(task) == 1506 assert result(task) == 1506
@@ -82,7 +82,7 @@ def test_sentinel():
def test_cluster(broker): def test_cluster(broker):
broker.list_key = 'cluster_test:q' broker.list_key = 'cluster_test:q'
broker.delete_queue() broker.delete_queue()
task = async('django_q.tests.tasks.count_letters', DEFAULT_WORDLIST, broker=broker) task = enqueue('django_q.tests.tasks.count_letters', DEFAULT_WORDLIST, broker=broker)
assert broker.queue_size() == 1 assert broker.queue_size() == 1
task_queue = Queue() task_queue = Queue()
assert task_queue.qsize() == 0 assert task_queue.qsize() == 0
@@ -109,32 +109,32 @@ def test_cluster(broker):
@pytest.mark.django_db @pytest.mark.django_db
def test_async(broker, admin_user): def test_enqueue(broker, admin_user):
broker.list_key = 'cluster_test:q' broker.list_key = 'cluster_test:q'
broker.delete_queue() broker.delete_queue()
a = async('django_q.tests.tasks.count_letters', DEFAULT_WORDLIST, hook='django_q.tests.test_cluster.assert_result', a = enqueue('django_q.tests.tasks.count_letters', DEFAULT_WORDLIST, hook='django_q.tests.test_cluster.assert_result',
broker=broker) broker=broker)
b = async('django_q.tests.tasks.count_letters2', WordClass(), hook='django_q.tests.test_cluster.assert_result', b = enqueue('django_q.tests.tasks.count_letters2', WordClass(), hook='django_q.tests.test_cluster.assert_result',
broker=broker) broker=broker)
# unknown argument # unknown argument
c = async('django_q.tests.tasks.count_letters', DEFAULT_WORDLIST, 'oneargumentoomany', c = enqueue('django_q.tests.tasks.count_letters', DEFAULT_WORDLIST, 'oneargumentoomany',
hook='django_q.tests.test_cluster.assert_bad_result', broker=broker) hook='django_q.tests.test_cluster.assert_bad_result', broker=broker)
# unknown function # unknown function
d = async('django_q.tests.tasks.does_not_exist', WordClass(), hook='django_q.tests.test_cluster.assert_bad_result', d = enqueue('django_q.tests.tasks.does_not_exist', WordClass(), hook='django_q.tests.test_cluster.assert_bad_result',
broker=broker) broker=broker)
# function without result # function without result
e = async('django_q.tests.tasks.countdown', 100000, broker=broker) e = enqueue('django_q.tests.tasks.countdown', 100000, broker=broker)
# function as instance # function as instance
f = async(multiply, 753, 2, hook=assert_result, broker=broker) f = enqueue(multiply, 753, 2, hook=assert_result, broker=broker)
# model as argument # model as argument
g = async('django_q.tests.tasks.get_task_name', Task(name='John'), broker=broker) g = enqueue('django_q.tests.tasks.get_task_name', Task(name='John'), broker=broker)
# args,kwargs, group and broken hook # args,kwargs, group and broken hook
h = async('django_q.tests.tasks.word_multiply', 2, word='django', hook='fail.me', broker=broker) h = enqueue('django_q.tests.tasks.word_multiply', 2, word='django', hook='fail.me', broker=broker)
# args unpickle test # args unpickle test
j = async('django_q.tests.tasks.get_user_id', admin_user, broker=broker, group='test_j') j = enqueue('django_q.tests.tasks.get_user_id', admin_user, broker=broker, group='test_j')
# q_options and save opt_out test # q_options and save opt_out test
k = async('django_q.tests.tasks.get_user_id', admin_user, k = enqueue('django_q.tests.tasks.get_user_id', admin_user,
q_options={'broker': broker, 'group': 'test_k', 'save': False, 'timeout': 90}) q_options={'broker': broker, 'group': 'test_k', 'save': False, 'timeout': 90})
# check if everything has a task id # check if everything has a task id
assert isinstance(a, str) assert isinstance(a, str)
assert isinstance(b, str) assert isinstance(b, str)
@@ -249,7 +249,7 @@ def test_timeout(broker):
# set up the Sentinel # set up the Sentinel
broker.list_key = 'timeout_test:q' broker.list_key = 'timeout_test:q'
broker.purge_queue() broker.purge_queue()
async('django_q.tests.tasks.count_forever', broker=broker) enqueue('django_q.tests.tasks.count_forever', broker=broker)
start_event = Event() start_event = Event()
stop_event = Event() stop_event = Event()
# Set a timer to stop the Sentinel # Set a timer to stop the Sentinel
@@ -265,7 +265,7 @@ def test_timeout(broker):
def test_timeout_override(broker): def test_timeout_override(broker):
# set up the Sentinel # set up the Sentinel
broker.list_key = 'timeout_override_test:q' broker.list_key = 'timeout_override_test:q'
async('django_q.tests.tasks.count_forever', broker=broker, timeout=1) enqueue('django_q.tests.tasks.count_forever', broker=broker, timeout=1)
start_event = Event() start_event = Event()
stop_event = Event() stop_event = Event()
# Set a timer to stop the Sentinel # Set a timer to stop the Sentinel
@@ -281,9 +281,9 @@ def test_timeout_override(broker):
def test_recycle(broker, monkeypatch): def test_recycle(broker, monkeypatch):
# set up the Sentinel # set up the Sentinel
broker.list_key = 'test_recycle_test:q' broker.list_key = 'test_recycle_test:q'
async('django_q.tests.tasks.multiply', 2, 2, broker=broker) enqueue('django_q.tests.tasks.multiply', 2, 2, broker=broker)
async('django_q.tests.tasks.multiply', 2, 2, broker=broker) enqueue('django_q.tests.tasks.multiply', 2, 2, broker=broker)
async('django_q.tests.tasks.multiply', 2, 2, broker=broker) enqueue('django_q.tests.tasks.multiply', 2, 2, broker=broker)
start_event = Event() start_event = Event()
stop_event = Event() stop_event = Event()
# override settings # override settings
@@ -295,8 +295,8 @@ def test_recycle(broker, monkeypatch):
assert start_event.is_set() assert start_event.is_set()
assert s.status() == Conf.STOPPED assert s.status() == Conf.STOPPED
assert s.reincarnations == 1 assert s.reincarnations == 1
async('django_q.tests.tasks.multiply', 2, 2, broker=broker) enqueue('django_q.tests.tasks.multiply', 2, 2, broker=broker)
async('django_q.tests.tasks.multiply', 2, 2, broker=broker) enqueue('django_q.tests.tasks.multiply', 2, 2, broker=broker)
task_queue = Queue() task_queue = Queue()
result_queue = Queue() result_queue = Queue()
# push two tasks # push two tasks
@@ -318,7 +318,7 @@ def test_recycle(broker, monkeypatch):
@pytest.mark.django_db @pytest.mark.django_db
def test_bad_secret(broker, monkeypatch): def test_bad_secret(broker, monkeypatch):
broker.list_key = 'test_bad_secret:q' broker.list_key = 'test_bad_secret:q'
async('math.copysign', 1, -1, broker=broker) enqueue('math.copysign', 1, -1, broker=broker)
stop_event = Event() stop_event = Event()
stop_event.set() stop_event.set()
start_event = Event() start_event = Event()
+2 -2
View File
@@ -1,6 +1,6 @@
import pytest import pytest
from django_q.tasks import async from django_q.tasks import enqueue
from django_q.brokers import get_broker from django_q.brokers import get_broker
from django_q.cluster import Cluster from django_q.cluster import Cluster
from django_q.compat import range from django_q.compat import range
@@ -46,4 +46,4 @@ def test_info():
def do_sync(): def do_sync():
async('django_q.tests.tasks.countdown', 1, sync=True, save=True) enqueue('django_q.tests.tasks.countdown', 1, sync=True, save=True)
+11 -11
View File
@@ -2,17 +2,17 @@
Chains Chains
====== ======
Sometimes you want to run tasks sequentially. For that you can use the :func:`async_chain` function: Sometimes you want to run tasks sequentially. For that you can use the :func:`enqueue_chain` function:
.. code-block:: python .. code-block:: python
# Async a chain of tasks # enqueue a chain of tasks
from django_q.tasks import async_chain, result_group from django_q.tasks import enqueue_chain, result_group
# the chain must be in the format # the chain must be in the format
# [(func,(args),{kwargs}),(func,(args),{kwargs}),..] # [(func,(args),{kwargs}),(func,(args),{kwargs}),..]
group_id = async_chain([('math.copysign', (1, -1)), group_id = enqueue_chain([('math.copysign', (1, -1)),
('math.floor', (1,))]) ('math.floor', (1,))])
# get group result # get group result
result_group(group_id, count=2) result_group(group_id, count=2)
@@ -21,7 +21,7 @@ A slightly more convenient way is to use a :class:`Chain` instance:
.. code-block:: python .. code-block:: python
# Chain async # Chain enqueue
from django_q.tasks import Chain from django_q.tasks import Chain
# create a chain that uses the cache backend # create a chain that uses the cache backend
@@ -41,9 +41,9 @@ A slightly more convenient way is to use a :class:`Chain` instance:
Reference Reference
--------- ---------
.. py:function:: async_chain(chain, group=None, cached=Conf.CACHED, sync=Conf.SYNC, broker=None) .. py:function:: enqueue_chain(chain, group=None, cached=Conf.CACHED, sync=Conf.SYNC, broker=None)
Async a chain of tasks. See also the :class:`Chain` class. enqueue a chain of tasks. See also the :class:`Chain` class.
:param list chain: a list of tasks in the format [(func,(args),{kwargs}), (func,(args),{kwargs})] :param list chain: a list of tasks in the format [(func,(args),{kwargs}), (func,(args),{kwargs})]
:param str group: an optional group name. :param str group: an optional group name.
@@ -52,7 +52,7 @@ Reference
.. py:class:: Chain(chain=None, group=None, cached=Conf.CACHED, sync=Conf.SYNC) .. py:class:: Chain(chain=None, group=None, cached=Conf.CACHED, sync=Conf.SYNC)
A sequential chain of tasks. Acts as a convenient wrapper for :func:`async_chain` A sequential chain of tasks. Acts as a convenient wrapper for :func:`enqueue_chain`
You can pass the task chain at construction or you can append individual tasks before running them. You can pass the task chain at construction or you can append individual tasks before running them.
:param list chain: a list of task in the format [(func,(args),{kwargs}), (func,(args),{kwargs})] :param list chain: a list of task in the format [(func,(args),{kwargs}), (func,(args),{kwargs})]
@@ -63,7 +63,7 @@ Reference
.. py:method:: append(func, *args, **kwargs) .. py:method:: append(func, *args, **kwargs)
Append a task to the chain. Takes the same arguments as :func:`async` Append a task to the chain. Takes the same arguments as :func:`enqueue`
:return: the current number of tasks in the chain :return: the current number of tasks in the chain
:rtype: int :rtype: int
@@ -102,4 +102,4 @@ Reference
get the length of the chain get the length of the chain
:return int: length of the chain :return int: length of the chain
+2 -2
View File
@@ -64,7 +64,7 @@ Set this to something that makes sense for your project. Can be overridden for i
ack_failures ack_failures
~~~~~~~~~~~~ ~~~~~~~~~~~~
When set to ``True``, also acknowledge unsuccessful tasks. This causes failed tasks to be considered as successful deliveries, thereby removing them from the task queue. Can also be set per-task by passing the ``ack_failure`` option to :func:`async`. Defaults to ``False``. When set to ``True``, also acknowledge unsuccessful tasks. This causes failed tasks to be considered as successful deliveries, thereby removing them from the task queue. Can also be set per-task by passing the ``ack_failure`` option to :func:`enqueue`. Defaults to ``False``.
.. _retry: .. _retry:
@@ -101,7 +101,7 @@ Guard loop sleep in seconds, must be greater than 0 and less than 60.
sync sync
~~~~ ~~~~
When set to ``True`` this configuration option forces all :func:`async` calls to be run with ``sync=True``. When set to ``True`` this configuration option forces all :func:`enqueue` calls to be run with ``sync=True``.
Effectively making everything synchronous. Useful for testing. Defaults to ``False``. Effectively making everything synchronous. Useful for testing. Defaults to ``False``.
.. _queue_limit: .. _queue_limit:
+44 -44
View File
@@ -12,18 +12,18 @@ Sending an email can take a while so why not queue it:
# Welcome mail with follow up example # Welcome mail with follow up example
from datetime import timedelta from datetime import timedelta
from django.utils import timezone from django.utils import timezone
from django_q.tasks import async, schedule from django_q.tasks import enqueue, schedule
from django_q.models import Schedule from django_q.models import Schedule
def welcome_mail(user): def welcome_mail(user):
msg = 'Welcome to our website' msg = 'Welcome to our website'
# send this message right away # send this message right away
async('django.core.mail.send_mail', enqueue('django.core.mail.send_mail',
'Welcome', 'Welcome',
msg, msg,
'from@example.com', 'from@example.com',
[user.email]) [user.email])
# and this follow up email in one hour # and this follow up email in one hour
msg = 'Here are some tips to get you started...' msg = 'Here are some tips to get you started...'
schedule('django.core.mail.send_mail', schedule('django.core.mail.send_mail',
@@ -51,7 +51,7 @@ A good place to use async tasks are Django's model signals. You don't want to de
from django.contrib.auth.models import User from django.contrib.auth.models import User
from django.db.models.signals import pre_save from django.db.models.signals import pre_save
from django.dispatch import receiver from django.dispatch import receiver
from django_q.tasks import async from django_q.tasks import enqueue
# set up the pre_save signal for our user # set up the pre_save signal for our user
@receiver(pre_save, sender=User) @receiver(pre_save, sender=User)
@@ -64,7 +64,7 @@ A good place to use async tasks are Django's model signals. You don't want to de
# has his email changed? # has his email changed?
if not user.email == instance.email: if not user.email == instance.email:
# tell everyone # tell everyone
async('tasks.inform_everyone', instance) enqueue('tasks.inform_everyone', instance)
The task will send a message to everyone else informing them that the users email address has changed. Note that this adds almost no overhead to the save action: The task will send a message to everyone else informing them that the users email address has changed. Note that this adds almost no overhead to the save action:
@@ -87,8 +87,8 @@ The task will send a message to everyone else informing them that the users emai
for u in User.objects.exclude(pk=user.pk): for u in User.objects.exclude(pk=user.pk):
msg = 'Dear {}, {} has a new email address: {}' msg = 'Dear {}, {} has a new email address: {}'
msg = msg.format(u.username, user.username, user.email) msg = msg.format(u.username, user.username, user.email)
async('django.core.mail.send_mail', enqueue('django.core.mail.send_mail',
'New email', msg, 'from@example.com', [u.email]) 'New email', msg, 'from@example.com', [u.email])
Of course you can do other things beside sending emails. These are just generic examples. You can use signals with async to update fields in other objects too. Of course you can do other things beside sending emails. These are just generic examples. You can use signals with async to update fields in other objects too.
@@ -104,19 +104,19 @@ In this example the user requests a report and we let the cluster do the generat
.. code-block:: python .. code-block:: python
# Report generation with hook example # Report generation with hook example
from django_q.tasks import async from django_q.tasks import enqueue
# views.py # views.py
# user requests a report. # user requests a report.
def create_report(request): def create_report(request):
async('tasks.create_html_report', enqueue('tasks.create_html_report',
request.user, request.user,
hook='tasks.email_report') hook='tasks.email_report')
.. code-block:: python .. code-block:: python
# tasks.py # tasks.py
from django_q.tasks import async from django_q.tasks import enqueue
# report generator # report generator
def create_html_report(user): def create_html_report(user):
@@ -127,16 +127,16 @@ In this example the user requests a report and we let the cluster do the generat
def email_report(task): def email_report(task):
if task.success: if task.success:
# Email the report # Email the report
async('django.core.mail.send_mail', enqueue('django.core.mail.send_mail',
'The report you requested', 'The report you requested',
task.result, task.result,
'from@example.com', 'from@example.com',
task.args[0].email) task.args[0].email)
else: else:
# Tell the admins something went wrong # Tell the admins something went wrong
async('django.core.mail.mail_admins', enqueue('django.core.mail.mail_admins',
'Report generation failed', 'Report generation failed',
task.result) task.result)
The hook is practical here, because it allows us to detach the sending task from the report generation function and to report on possible failures. The hook is practical here, because it allows us to detach the sending task from the report generation function and to report on possible failures.
@@ -152,12 +152,12 @@ here's an example of how you can have Django Q take care of your indexes in real
from .models import Document from .models import Document
from django.db.models.signals import post_save from django.db.models.signals import post_save
from django.dispatch import receiver from django.dispatch import receiver
from django_q.tasks import async from django_q.tasks import enqueue
# hook up the post save handler # hook up the post save handler
@receiver(post_save, sender=Document) @receiver(post_save, sender=Document)
def document_changed(sender, instance, **kwargs): def document_changed(sender, instance, **kwargs):
async('tasks.index_object', sender, instance, save=False) enqueue('tasks.index_object', sender, instance, save=False)
# turn off result saving to not flood your database # turn off result saving to not flood your database
.. code-block:: python .. code-block:: python
@@ -177,7 +177,7 @@ here's an example of how you can have Django Q take care of your indexes in real
index.update_object(instance, using=backend) index.update_object(instance, using=backend)
Now every time a Document is saved, your indexes will be updated without causing a delay in your save action. Now every time a Document is saved, your indexes will be updated without causing a delay in your save action.
You could expand this to dealing with deletes, by adding a ``post_delete`` signal and calling ``index.remove_object`` in the async function. You could expand this to dealing with deletes, by adding a ``post_delete`` signal and calling ``index.remove_object`` in the enqueue function.
.. _shell: .. _shell:
@@ -187,13 +187,13 @@ You can execute or schedule shell commands using Pythons :mod:`subprocess` modul
.. code-block:: python .. code-block:: python
from django_q.tasks import async, result from django_q.tasks import enqueue, result
# make a backup copy of setup.py # make a backup copy of setup.py
async('subprocess.call', ['cp', 'setup.py', 'setup.py.bak']) enqueue('subprocess.call', ['cp', 'setup.py', 'setup.py.bak'])
# call ls -l and dump the output # call ls -l and dump the output
task_id=async('subprocess.check_output', ['ls', '-l']) task_id=enqueue('subprocess.check_output', ['ls', '-l'])
# get the result # get the result
dir_list = result(task_id) dir_list = result(task_id)
@@ -202,10 +202,10 @@ In Python 3.5 the subprocess module has changed quite a bit and returns a :class
.. code-block:: python .. code-block:: python
from django_q.tasks import async, result from django_q.tasks import enqueue, result
# make a backup copy of setup.py # make a backup copy of setup.py
tid = async('subprocess.run', ['cp', 'setup.py', 'setup.py.bak']) tid = enqueue('subprocess.run', ['cp', 'setup.py', 'setup.py.bak'])
# get the result # get the result
r=result(tid, 500) r=result(tid, 500)
@@ -220,22 +220,22 @@ In Python 3.5 the subprocess module has changed quite a bit and returns a :class
from subprocess import PIPE from subprocess import PIPE
# call ls -l and pipe the output # call ls -l and pipe the output
tid = async('subprocess.run', ['ls', '-l'], stdout=PIPE) tid = enqueue('subprocess.run', ['ls', '-l'], stdout=PIPE)
# get the result # get the result
res = result(tid, 500) res = result(tid, 500)
# print the output # print the output
print(res.stdout) print(res.stdout)
Instead of :func:`async` you can of course also use :func:`schedule` to schedule commands. Instead of :func:`enqueue` you can of course also use :func:`schedule` to schedule commands.
For regular Django management commands, it is easier to call them directly: For regular Django management commands, it is easier to call them directly:
.. code-block:: python .. code-block:: python
from django_q.tasks import async, schedule from django_q.tasks import enqueue, schedule
async('django.core.management.call_command','clearsessions') enqueue('django.core.management.call_command','clearsessions')
# or clear those sessions every hour # or clear those sessions every hour
@@ -255,7 +255,7 @@ Adapted from `Sebastian Raschka's blog <http://sebastianraschka.com/Articles/201
# Group example with Parzen-window estimation # Group example with Parzen-window estimation
import numpy import numpy
from django_q.tasks import async, result_group, delete_group from django_q.tasks import enqueue, result_group, delete_group
# the estimation function # the estimation function
def parzen_estimation(x_samples, point_x, h): def parzen_estimation(x_samples, point_x, h):
@@ -270,7 +270,7 @@ Adapted from `Sebastian Raschka's blog <http://sebastianraschka.com/Articles/201
return h, (k_n / len(x_samples)) / (h ** point_x.shape[1]) return h, (k_n / len(x_samples)) / (h ** point_x.shape[1])
# create 100 calculations and return the collated result # create 100 calculations and return the collated result
def parzen_async(): def parzen_enqueue():
# clear the previous results # clear the previous results
delete_group('parzen', cached=True) delete_group('parzen', cached=True)
mu_vec = numpy.array([0, 0]) mu_vec = numpy.array([0, 0])
@@ -279,10 +279,10 @@ Adapted from `Sebastian Raschka's blog <http://sebastianraschka.com/Articles/201
multivariate_normal(mu_vec, cov_mat, 10000) multivariate_normal(mu_vec, cov_mat, 10000)
widths = numpy.linspace(1.0, 1.2, 100) widths = numpy.linspace(1.0, 1.2, 100)
x = numpy.array([[0], [0]]) x = numpy.array([[0], [0]])
# async them with a group label to the cache backend # enqueue them with a group label to the cache backend
for w in widths: for w in widths:
async(parzen_estimation, sample, x, w, enqueue(parzen_estimation, sample, x, w,
group='parzen', cached=True) group='parzen', cached=True)
# return after 100 results # return after 100 results
return result_group('parzen', count=100, cached=True) return result_group('parzen', count=100, cached=True)
@@ -290,21 +290,21 @@ Adapted from `Sebastian Raschka's blog <http://sebastianraschka.com/Articles/201
Django Q is not optimized for distributed computing, but this example will give you an idea of what you can do with task :doc:`group`. Django Q is not optimized for distributed computing, but this example will give you an idea of what you can do with task :doc:`group`.
Alternatively the ``parzen_async()`` function can also be written with :func:`async_iter`, which automatically utilizes the cache backend and groups to return a single result from an iterable: Alternatively the ``parzen_enqueue()`` function can also be written with :func:`enqueue_iter`, which automatically utilizes the cache backend and groups to return a single result from an iterable:
.. code-block:: python .. code-block:: python
# create 100 calculations and return the collated result # create 100 calculations and return the collated result
def parzen_async(): def parzen_enqueue():
mu_vec = numpy.array([0, 0]) mu_vec = numpy.array([0, 0])
cov_mat = numpy.array([[1, 0], [0, 1]]) cov_mat = numpy.array([[1, 0], [0, 1]])
sample = numpy.random. \ sample = numpy.random. \
multivariate_normal(mu_vec, cov_mat, 10000) multivariate_normal(mu_vec, cov_mat, 10000)
widths = numpy.linspace(1.0, 1.2, 100) widths = numpy.linspace(1.0, 1.2, 100)
x = numpy.array([[0], [0]]) x = numpy.array([[0], [0]])
# async them with async iterable # enqueue them with enqueue iterable
args = [(sample, x, w) for w in widths] args = [(sample, x, w) for w in widths]
result_id = async_iter(parzen_estimation, args, cached=True) result_id = enqueue_iter(parzen_estimation, args, cached=True)
# return the cached result or timeout after 10 seconds # return the cached result or timeout after 10 seconds
return result(result_id, wait=10000, cached=True) return result(result_id, wait=10000, cached=True)
+7 -7
View File
@@ -2,15 +2,15 @@
Groups Groups
====== ======
You can group together results by passing :func:`async` the optional ``group`` keyword: You can group together results by passing :func:`enqueue` the optional ``group`` keyword:
.. code-block:: python .. code-block:: python
# result group example # result group example
from django_q.tasks import async, result_group from django_q.tasks import enqueue, result_group
for i in range(4): for i in range(4):
async('math.modf', i, group='modf') enqueue('math.modf', i, group='modf')
# wait until the group has 4 results # wait until the group has 4 results
result = result_group('modf', count=4) result = result_group('modf', count=4)
@@ -66,14 +66,14 @@ You can also access group functions from a task result instance:
task.group_delete() task.group_delete()
print('Deleted group {}'.format(task.group)) print('Deleted group {}'.format(task.group))
or call them directly on :class:`Async` object: or call them directly on :class:`AsyncTask` object:
.. code-block:: python .. code-block:: python
from django_q.tasks import Async from django_q.tasks import enqueue
# add a task to the math group and run it cached # add a task to the math group and run it cached
a = Async('math.floor', 2.5, group='math', cached=True) a = enqueue('math.floor', 2.5, group='math', cached=True)
# wait until this tasks group has 10 results # wait until this tasks group has 10 results
result = a.result_group(count=10) result = a.result_group(count=10)
@@ -122,4 +122,4 @@ Reference
:param bool tasks: also deletes the associated tasks if ``True`` :param bool tasks: also deletes the associated tasks if ``True``
:param bool cached: run this against the cache backend. :param bool cached: run this against the cache backend.
:returns: the numbers of tasks affected :returns: the numbers of tasks affected
:rtype: int :rtype: int
+7 -7
View File
@@ -2,16 +2,16 @@
Iterable Iterable
======== ========
If you have an iterable object with arguments for a function, you can use :func:`async_iter` to async them with a single command:: If you have an iterable object with arguments for a function, you can use :func:`enqueue_iter` to async them with a single command::
# Async Iterable example # Async Iterable example
from django_q.tasks import async_iter, result from django_q.tasks import enqueue_iter, result
# set up a list of arguments for math.floor # set up a list of arguments for math.floor
iter = [i for i in range(100)] iter = [i for i in range(100)]
# async iter them # enqueue iter them
id=async_iter('math.floor',iter) id=enqueue_iter('math.floor',iter)
# wait for the collated result for 1 second # wait for the collated result for 1 second
result_list = result(id, wait=1000) result_list = result(id, wait=1000)
@@ -45,10 +45,10 @@ You can also use an :class:`Iter` instance which can sometimes be more convenien
Reference Reference
--------- ---------
.. py:function:: async_iter(func, args_iter,**kwargs) .. py:function:: enqueue_iter(func, args_iter,**kwargs)
Runs iterable arguments against the cache backend and returns a single collated result. Runs iterable arguments against the cache backend and returns a single collated result.
Accepts the same options as :func:`async` except ``hook``. See also the :class:`Iter` class. Accepts the same options as :func:`enqueue` except ``hook``. See also the :class:`Iter` class.
:param object func: The task function to execute :param object func: The task function to execute
:param args: An iterable containing arguments for the task function :param args: An iterable containing arguments for the task function
@@ -58,7 +58,7 @@ Reference
.. py:class:: Iter(func=None, args=None, kwargs=None, cached=Conf.CACHED, sync=Conf.SYNC, broker=None) .. py:class:: Iter(func=None, args=None, kwargs=None, cached=Conf.CACHED, sync=Conf.SYNC, broker=None)
An async task with iterable arguments. Serves as a convenient wrapper for :func:`async_iter` An async task with iterable arguments. Serves as a convenient wrapper for :func:`enqueue_iter`
You can pass the iterable arguments at construction or you can append individual argument tuples. You can pass the iterable arguments at construction or you can append individual argument tuples.
:param func: the function to execute :param func: the function to execute
+2 -2
View File
@@ -27,7 +27,7 @@ You can manage them through the :ref:`admin_page` or directly from your code wit
schedule_type=Schedule.DAILY schedule_type=Schedule.DAILY
) )
# In case you want to use async options # In case you want to use q_options
schedule('math.sqrt', schedule('math.sqrt',
9, 9,
hook='hooks.print_result', hook='hooks.print_result',
@@ -103,7 +103,7 @@ Reference
:param int minutes: Number of minutes for the Minutes type. :param int minutes: Number of minutes for the Minutes type.
:param int repeats: Number of times to repeat schedule. -1=Always, 0=Never, n =n. :param int repeats: Number of times to repeat schedule. -1=Always, 0=Never, n =n.
:param datetime next_run: Next or first scheduled execution datetime. :param datetime next_run: Next or first scheduled execution datetime.
:param dict q_options: async options to use for this schedule :param dict q_options: options passed to enqueue for this schedule
:param kwargs: optional keyword arguments for the scheduled function. :param kwargs: optional keyword arguments for the scheduled function.
.. class:: Schedule .. class:: Schedule
+41 -41
View File
@@ -4,22 +4,22 @@ Tasks
.. _async: .. _async:
async() enqueue()
------- ---------
Use :func:`async` from your code to quickly offload tasks to the :class:`Cluster`: Use :func:`enqueue` from your code to quickly offload tasks to the :class:`Cluster`:
.. code:: python .. code:: python
from django_q.tasks import async, result from django_q.tasks import enqueue, result
# create the task # create the task
async('math.copysign', 2, -2) enqueue('math.copysign', 2, -2)
# or with import and storing the id # or with import and storing the id
import math.copysign import math.copysign
task_id = async(copysign, 2, -2) task_id = enqueue(copysign, 2, -2)
# get the result # get the result
task_result = result(task_id) task_result = result(task_id)
@@ -30,13 +30,13 @@ Use :func:`async` from your code to quickly offload tasks to the :class:`Cluster
# but in most cases you will want to use a hook: # but in most cases you will want to use a hook:
async('math.modf', 2.5, hook='hooks.print_result') enqueue('math.modf', 2.5, hook='hooks.print_result')
# hooks.py # hooks.py
def print_result(task): def print_result(task):
print(task.result) print(task.result)
:func:`async` can take the following optional keyword arguments: :func:`enqueue` can take the following optional keyword arguments:
hook hook
"""" """"
@@ -84,13 +84,13 @@ None of the option keywords get passed on to the task function.
As an alternative you can also put them in As an alternative you can also put them in
a single keyword dict named ``q_options``. This enables you to use these keywords for your function call:: a single keyword dict named ``q_options``. This enables you to use these keywords for your function call::
# Async options in a dict # Enqueue options in a dict
opts = {'hook': 'hooks.print_result', opts = {'hook': 'hooks.print_result',
'group': 'math', 'group': 'math',
'timeout': 30} 'timeout': 30}
async('math.modf', 2.5, q_options=opts) enqueue('math.modf', 2.5, q_options=opts)
Please note that this will override any other option keywords. Please note that this will override any other option keywords.
@@ -99,18 +99,18 @@ Please note that this will override any other option keywords.
or you need to configure Django Q to run in synchronous mode for testing using the :ref:`sync` option. or you need to configure Django Q to run in synchronous mode for testing using the :ref:`sync` option.
Async AsyncTask
----- ---------
Optionally you can use the :class:`Async` class to instantiate a task and keep everything in a single object.: Optionally you can use the :class:`AsyncTask` class to instantiate a task and keep everything in a single object.:
.. code-block:: python .. code-block:: python
# Async class instance example # AsyncTask class instance example
from django_q.tasks import Async from django_q.tasks import AsyncTask
# instantiate an async task # instantiate an async task
a = Async('math.floor', 1.5, group='math') a = AsyncTask('math.floor', 1.5, group='math')
# you can set or change keywords afterwards # you can set or change keywords afterwards
a.cached = True a.cached = True
@@ -136,7 +136,7 @@ Optionally you can use the :class:`Async` class to instantiate a task and keep e
1 1
2 2
Once you change any of the parameters of the task after it has run, the result is invalidated and you will have to :func:`Async.run` it again to retrieve a new result. Once you change any of the parameters of the task after it has run, the result is invalidated and you will have to :func:`AsyncTask.run` it again to retrieve a new result.
Cached operations Cached operations
----------------- -----------------
@@ -150,10 +150,10 @@ You can also opt to set a manual timeout on the results, by setting e.g. ``cache
This works both globally or on individual async executions.:: This works both globally or on individual async executions.::
# simple cached example # simple cached example
from django_q.tasks import async, result from django_q.tasks import enqueue, result
# cache the result for 10 seconds # cache the result for 10 seconds
id = async('math.floor', 100, cached=10) id = enqueue('math.floor', 100, cached=10)
# wait max 50ms for the result to appear in the cache # wait max 50ms for the result to appear in the cache
result(id, wait=50, cached=True) result(id, wait=50, cached=True)
@@ -169,35 +169,35 @@ As you can see you can easily turn a cached result into a permanent database res
This also works for group actions:: This also works for group actions::
# cached group example # cached group example
from django_q.tasks import async, result_group from django_q.tasks import enqueue, result_group
from django_q.brokers import get_broker from django_q.brokers import get_broker
# set up a broker instance for better performance # set up a broker instance for better performance
broker = get_broker() broker = get_broker()
# async a hundred functions under a group label # enqueue a hundred functions under a group label
for i in range(100): for i in range(100):
async('math.frexp', enqueue('math.frexp',
i, i,
group='frexp', group='frexp',
cached=True, cached=True,
broker=broker) broker=broker)
# wait max 50ms for one hundred results to return # wait max 50ms for one hundred results to return
result_group('frexp', wait=50, count=100, cached=True) result_group('frexp', wait=50, count=100, cached=True)
If you don't need hooks, that exact same result can be achieved by using the more convenient :func:`async_iter`. If you don't need hooks, that exact same result can be achieved by using the more convenient :func:`enqueue_iter`.
Synchronous testing Synchronous testing
------------------- -------------------
:func:`async` can be instructed to execute a task immediately by setting the optional keyword ``sync=True``. :func:`enqueue` can be instructed to execute a task immediately by setting the optional keyword ``sync=True``.
The task will then be injected straight into a worker and the result saved by a monitor instance:: The task will then be injected straight into a worker and the result saved by a monitor instance::
from django_q.tasks import async, fetch from django_q.tasks import enqueue, fetch
# create a synchronous task # create a synchronous task
task_id = async('my.buggy.code', sync=True) task_id = enqueue('my.buggy.code', sync=True)
# the task will then be available immediately # the task will then be available immediately
task = fetch(task_id) task = fetch(task_id)
@@ -210,24 +210,24 @@ The task will then be injected straight into a worker and the result saved by a
An error occurred: ImportError("No module named 'my'",) An error occurred: ImportError("No module named 'my'",)
Note that :func:`async` will block until the task is executed and saved. This feature bypasses the broker and is intended for debugging and development. Note that :func:`enqueue` will block until the task is executed and saved. This feature bypasses the broker and is intended for debugging and development.
Instead of setting ``sync`` on each individual ``async`` you can also configure :ref:`sync` as a global override. Instead of setting ``sync`` on each individual ``enqueue`` you can also configure :ref:`sync` as a global override.
Connection pooling Connection pooling
------------------ ------------------
Django Q tries to pass broker instances around its parts as much as possible to save you from running out of connections. Django Q tries to pass broker instances around its parts as much as possible to save you from running out of connections.
When you are making individual calls to :func:`async` a lot though, it can help to set up a broker to reuse for :func:`async`: When you are making individual calls to :func:`enqueue` a lot though, it can help to set up a broker to reuse for :func:`enqueue`:
.. code:: python .. code:: python
# broker connection economy example # broker connection economy example
from django_q.tasks import async from django_q.tasks import enqueue
from django_q.brokers import get_broker from django_q.brokers import get_broker
broker = get_broker() broker = get_broker()
for i in range(50): for i in range(50):
async('math.modf', 2.5, broker=broker) enqueue('math.modf', 2.5, broker=broker)
.. tip:: .. tip::
@@ -237,7 +237,7 @@ When you are making individual calls to :func:`async` a lot though, it can help
Reference Reference
--------- ---------
.. py:function:: async(func, *args, hook=None, group=None, timeout=None,\ .. py:function:: enqueue(func, *args, hook=None, group=None, timeout=None,\
save=None, sync=False, cached=False, broker=None, q_options=None, **kwargs) save=None, sync=False, cached=False, broker=None, q_options=None, **kwargs)
Puts a task in the cluster queue Puts a task in the cluster queue
@@ -249,7 +249,7 @@ Reference
:param int timeout: Overrides global cluster :ref:`timeout`. :param int timeout: Overrides global cluster :ref:`timeout`.
:param bool save: Overrides global save setting for this task. :param bool save: Overrides global save setting for this task.
:param bool ack_failure: Overrides the global :ref:`ack_failures` setting for this task. :param bool ack_failure: Overrides the global :ref:`ack_failures` setting for this task.
:param bool sync: If set to True, async will simulate a task execution :param bool sync: If set to True, enqueue will simulate a task execution
:param cached: Output the result to the cache backend. Bool or timeout in seconds :param cached: Output the result to the cache backend. Bool or timeout in seconds
:param broker: Optional broker connection from :func:`brokers.get_broker` :param broker: Optional broker connection from :func:`brokers.get_broker`
:param dict q_options: Options dict, overrides option keywords :param dict q_options: Options dict, overrides option keywords
@@ -408,13 +408,13 @@ Reference
A proxy model of :class:`Task` with the queryset filtered on :attr:`Task.success` is ``False``. A proxy model of :class:`Task` with the queryset filtered on :attr:`Task.success` is ``False``.
.. py:class:: Async(func, *args, **kwargs) .. py:class:: AsyncTask(func, *args, **kwargs)
A class wrapper for the :func:`async` function. A class wrapper for the :func:`enqueue` function.
:param object func: The task function to execute :param object func: The task function to execute
:param tuple args: The arguments for the task function :param tuple args: The arguments for the task function
:param dict kwargs: Keyword arguments for the task function, including async options :param dict kwargs: Keyword arguments for the task function, including enqueue options
.. py:attribute:: id .. py:attribute:: id
@@ -434,7 +434,7 @@ Reference
.. py:attribute:: kwargs .. py:attribute:: kwargs
Keyword arguments for the function. Can include any of the optional async keyword attributes directly or in a `q_options` dictionary. Keyword arguments for the function. Can include any of the optional enqueue keyword attributes directly or in a `q_options` dictionary.
.. py:attribute:: broker .. py:attribute:: broker