From f5197ac9120b8eb7b87d8edf33d5d995b266e296 Mon Sep 17 00:00:00 2001 From: Ilan Steemers Date: Wed, 28 Oct 2015 12:06:18 +0100 Subject: [PATCH] docs: added Async class Also moved Group, Iterable and Chain to their own pages because the Task page was getting too big. v0.7.11 --- docs/chain.rst | 105 +++++++++++++ docs/conf.py | 2 +- docs/examples.rst | 2 +- docs/group.rst | 125 +++++++++++++++ docs/index.rst | 3 + docs/iterable.rst | 99 ++++++++++++ docs/tasks.rst | 377 +++++++++++----------------------------------- 7 files changed, 423 insertions(+), 290 deletions(-) create mode 100644 docs/chain.rst create mode 100644 docs/group.rst create mode 100644 docs/iterable.rst diff --git a/docs/chain.rst b/docs/chain.rst new file mode 100644 index 0000000..11298bc --- /dev/null +++ b/docs/chain.rst @@ -0,0 +1,105 @@ +.. py:currentmodule:: django_q + +Chains +====== +Sometimes you want to run tasks sequentially. For that you can use the :func:`async_chain` function: + +.. code-block:: python + + # Async a chain of tasks + from django_q.tasks import async_chain, result_group + + # the chain must be in the format + # [(func,(args),{kwargs}),(func,(args),{kwargs}),..] + group_id = async_chain([('math.copysign', (1, -1)), + ('math.floor', (1,))]) + + # get group result + result_group(group_id, count=2) + +A slightly more convenient way is to use a :class:`Chain` instance: + +.. code-block:: python + + # Chain async + from django_q.tasks import Chain + + # create a chain that uses the cache backend + chain = Chain(cached=True) + + # add some tasks + chain.append('math.copysign', 1, -1) + chain.append('math.floor', 1) + + # run it + chain.run() + + print(chain.result()) +.. code-block:: python + + [-1.0, 1] + +Reference +--------- +.. py:function:: async_chain(chain, group=None, cached=Conf.CACHED, sync=Conf.SYNC, broker=None) + + Async 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 str group: an optional group name. + :param bool cached: run this against the cache backend + :param bool sync: execute this inline instead of asynchronous + +.. 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` + 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 str group: an optional group name. + :param bool cached: run this against the cache backend + :param bool sync: execute this inline instead of asynchronous + + + .. py:method:: append(func, *args, **kwargs) + + Append a task to the chain. Takes the same arguments as :func:`async` + + :return: the current number of tasks in the chain + :rtype: int + + + .. py:method:: run() + + Start queueing the chain to the worker cluster. + + :return: the chains group id + + + .. py:method:: result(wait=0) + + return the full list of results from the chain when it finishes. Blocks until timeout or result. + + :param int wait: how many milliseconds to wait for a result + :return: an unsorted list of results + + + .. py:method:: fetch(failures=True, wait=0) + + get the task result objects from the chain when it finishes. Blocks until timeout or result. + + :param failures: include failed tasks + :param int wait: how many milliseconds to wait for a result + :return: an unsorted list of task objects + + .. py:method:: current() + + get the index of the currently executing chain element + + :return int: current chain index + + .. py:method:: length() + + get the length of the chain + + :return int: length of the chain \ No newline at end of file diff --git a/docs/conf.py b/docs/conf.py index abb2402..7e30354 100644 --- a/docs/conf.py +++ b/docs/conf.py @@ -72,7 +72,7 @@ author = 'Ilan Steemers' # The short X.Y version. version = '0.7' # The full version, including alpha/beta/rc tags. -release = '0.7.9' +release = '0.7.11' # The language for content autogenerated by Sphinx. Refer to documentation # for a list of supported languages. diff --git a/docs/examples.rst b/docs/examples.rst index fb6da62..bc6158f 100644 --- a/docs/examples.rst +++ b/docs/examples.rst @@ -288,7 +288,7 @@ Adapted from `Sebastian Raschka's blog 100: + print(task.group_result()) + task.group_delete() + print('Deleted group {}'.format(task.group)) + +or call them directly on :class:`Async` object: + +.. code-block:: python + + from django_q.tasks import Async + + # add a task to the math group and run it cached + a = Async('math.floor', 2.5, group='math', cached=True) + + # wait until this tasks group has 10 results + result = a.result_group(count=10) + +Reference +--------- +.. py:function:: result_group(group_id, failures=False, wait=0, count=None, cached=False) + + Returns the results of a task group + + :param str group_id: the group identifier + :param bool failures: set this to ``True`` to include failed results + :param int wait: optional milliseconds to wait for a result or count. -1 for indefinite + :param int count: block until there are this many results in the group + :param bool cached: run this against the cache backend + :returns: a list of results + :rtype: list + +.. py:function:: fetch_group(group_id, failures=True, wait=0, count=None, cached=False) + + Returns a list of tasks in a group + + :param str group_id: the group identifier + :param bool failures: set this to ``False`` to exclude failed tasks + :param int wait: optional milliseconds to wait for a task or count. -1 for indefinite + :param int count: block until there are this many tasks in the group + :param bool cached: run this against the cache backend. + :returns: a list of :class:`Task` + :rtype: list + +.. py:function:: count_group(group_id, failures=False, cached=False) + + Counts the number of task results in a group. + + :param str group_id: the group identifier + :param bool failures: counts the number of failures if ``True`` + :param bool cached: run this against the cache backend. + :returns: the number of tasks or failures in a group + :rtype: int + +.. py:function:: delete_group(group_id, tasks=False, cached=False) + + Deletes a group label from the database. + + :param str group_id: the group identifier + :param bool tasks: also deletes the associated tasks if ``True`` + :param bool cached: run this against the cache backend. + :returns: the numbers of tasks affected + :rtype: int \ No newline at end of file diff --git a/docs/index.rst b/docs/index.rst index 481c874..2e1aa0d 100644 --- a/docs/index.rst +++ b/docs/index.rst @@ -35,6 +35,9 @@ Contents: Configuration Brokers Tasks + Groups + Iterable + Chains Schedules Cluster Monitor diff --git a/docs/iterable.rst b/docs/iterable.rst new file mode 100644 index 0000000..b904ba5 --- /dev/null +++ b/docs/iterable.rst @@ -0,0 +1,99 @@ +.. py:currentmodule:: django_q + +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:: + + # Async Iterable example + from django_q.tasks import async_iter, result + + # set up a list of arguments for math.floor + iter = [i for i in range(100)] + + # async iter them + id=async_iter('math.floor',iter) + + # wait for the collated result for 1 second + result_list = result(id, wait=1000) + +This will individually queue 100 tasks to the worker cluster, which will save their results in the cache backend for speed. +Once all the 100 results are in the cache, they are collated into a list and saved as a single result in the database. The cache results are then cleared. + +You can also use an :class:`Iter` instance which can sometimes be more convenient: + +.. code-block:: python + + from django_q.tasks import Iter + + i = Iter('math.copysign') + + # add some arguments + i.append(1, -1) + i.append(2, -1) + i.append(3, -1) + +Reference +--------- + +.. py:function:: async_iter(func, args_iter,**kwargs) + + 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. + + :param object func: The task function to execute + :param args: An iterable containing arguments for the task function + :param dict kwargs: Keyword arguments for the task function. Ignores ``hook``. + :returns: The uuid of the task + :rtype: str + +.. 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` + You can pass the iterable arguments at construction or you can append individual argument tuples. + + :param func: the function to execute + :param args: an iterable of arguments. + :param kwargs: the keyword arguments + :param bool cached: run this against the cache backend + :param bool sync: execute this inline instead of asynchronous + :param broker: optional broker instance + + + .. py:method:: append(*args) + + Append arguments to the iter set. Returns the current set count. + + :param args: the arguments for a single execution + :return: the current set count + :rtype: int + + + .. py:method:: run() + + Start queueing the tasks to the worker cluster. + + :return: the task result id + + + .. py:method:: result(wait=0) + + return the full list of results. + + :param int wait: how many milliseconds to wait for a result + :return: an unsorted list of results + + + .. py:method:: fetch(wait=0) + + get the task result objects. + + :param int wait: how many milliseconds to wait for a result + :return: an unsorted list of task objects + + + .. py:method:: length() + + get the length of the arguments list + + :return int: length of the argument list + diff --git a/docs/tasks.rst b/docs/tasks.rst index 6d337ef..2f2abe7 100644 --- a/docs/tasks.rst +++ b/docs/tasks.rst @@ -4,8 +4,8 @@ Tasks .. _async: -Async ------ +async() +------- Use :func:`async` from your code to quickly offload tasks to the :class:`Cluster`: @@ -44,7 +44,7 @@ The function to call after the task has been executed. This function gets passed group """"" -A group label. Check :ref:`groups` for group functions. +A group label. Check :doc:`group` for group functions. save """" @@ -90,157 +90,44 @@ Please not 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. +Async +----- -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:: - - # Async Iterable example - from django_q.tasks import async_iter, result - - # set up a list of arguments for math.floor - iter = [i for i in range(100)] - - # async iter them - id=async_iter('math.floor',iter) - - # wait for the collated result for 1 second - result_list = result(id, wait=1000) - -This will individually queue 100 tasks to the worker cluster, which will save their results in the cache backend for speed. -Once all the 100 results are in the cache, they are collated into a list and saved as a single result in the database. The cache results are then cleared. - -You can also use an :class:`Iter` instance which can sometimes be more convenient: +Optionally you can use the :class:`Async` class to instantiate a task and keep everything in a single object.: .. code-block:: python - from django_q.tasks import Iter + # Async class instance example + from django_q.tasks import Async - i = Iter('math.copysign') + # instantiate an async task + a = Async('math.floor', 1.5, group='math') - # add some arguments - i.append(1, -1) - i.append(2, -1) - i.append(3, -1) + # you can set or change keywords afterwards + a.cached = True # run it - i.run() + a.run() - # get the results - print(i.result()) + # wait indefinitely for the result and print it + print(a.result(wait=-1)) + + # change the args + a.args = (2.5,) + + # run it again + a.run() + + # wait max 10 seconds for the result and print it + + print(a.result(wait=10)) .. code-block:: python - [-1.0, -2.0, -3.0] + 1 + 2 -Needs the Django cache framework. - -.. _groups: - -Groups ------- -You can group together results by passing :func:`async` the optional ``group`` keyword: - -.. code-block:: python - - # result group example - from django_q.tasks import async, result_group - - for i in range(4): - async('math.modf', i, group='modf') - - # wait until the group has 4 results - result = result_group('modf', count=4) - print(result) - -.. code-block:: python - - [(0.0, 0.0), (0.0, 1.0), (0.0, 2.0), (0.0, 3.0)] - -Note that the same can be achieved much faster with :func:`async_iter` - -Take care to not limit your results database too much and call :func:`delete_group` before each run, unless you want your results to keep adding up. -Instead of :func:`result_group` you can also use :func:`fetch_group` to return a queryset of :class:`Task` objects.: - -.. code-block:: python - - # fetch group example - from django_q.tasks import fetch_group, count_group, result_group - - # count the number of failures - failure_count = count_group('modf', failures=True) - - # only use the successes - results = fetch_group('modf') - if failure_count: - results = results.exclude(success=False) - results = [task.result for task in successes] - - # this is the same as - results = fetch_group('modf', failures=False) - results = [task.result for task in successes] - - # and the same as - results = result_group('modf') # filters failures by default - - -Getting results by using :func:`result_group` is of course much faster than using :func:`fetch_group`, but it doesn't offer the benefits of Django's queryset functions. - -.. note:: - - Calling ``Queryset.values`` for the result on Django 1.7 or lower will return a list of encoded results. - If you can't upgrade to Django 1.8, use list comprehension or an iterator to return decoded results. - -You can also access group functions from a task result instance: - -.. code-block:: python - - from django_q.tasks import fetch - - task = fetch('winter-speaker-alpha-ceiling') - if task.group_count() > 100: - print(task.group_result()) - task.group_delete() - print('Deleted group {}'.format(task.group)) - -Chains ------- -Sometimes you want to run tasks sequentially. For that you can use the :func:`async_chain` function: - -.. code-block:: python - - # Async a chain of tasks - from django_q.tasks import async_chain, result_group - - # the chain must be in the format - # [(func,(args),{kwargs}),(func,(args),{kwargs}),..] - group_id = async_chain([('math.copysign', (1, -1)), - ('math.floor', (1,))]) - - # get group result - result_group(group_id, count=2) - -A slightly more convenient way is to use a :class:`Chain` instance: - -.. code-block:: python - - # Chain async - from django_q.tasks import Chain - - # create a chain that uses the cache backend - chain = Chain(cached=True) - - # add some tasks - chain.append('math.copysign', 1, -1) - chain.append('math.floor', 1) - - # run it - chain.run() - - print(chain.result()) -.. code-block:: python - - [-1.0, 1] +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. Cached operations ----------------- @@ -365,7 +252,7 @@ Reference Gets the result of a previously executed task :param str task_id: the uuid or name of the task - :param int wait: optional milliseconds to wait for a result + :param int wait: optional milliseconds to wait for a result. -1 for indefinite :param bool cached: run this against the cache backend. :returns: The result of the executed task @@ -374,7 +261,7 @@ Reference Returns a previously executed task :param str name: the uuid or name of the task - :param int wait: optional milliseconds to wait for a result + :param int wait: optional milliseconds to wait for a result. -1 for indefinite :param bool cached: run this against the cache backend. :returns: A task object :rtype: Task @@ -383,27 +270,6 @@ Reference Renamed from get_task -.. py:function:: async_iter(func, args_iter,**kwargs) - - 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. - - :param object func: The task function to execute - :param args: An iterable containing arguments for the task function - :param dict kwargs: Keyword arguments for the task function. Ignores ``hook``. - :returns: The uuid of the task - :rtype: str - - -.. py:function:: async_chain(chain, group=None, cached=Conf.CACHED, sync=Conf.SYNC, broker=None) - - Async 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 str group: an optional group name. - :param bool cached: run this against the cache backend - :param bool sync: execute this inline instead of asynchronous - .. py:function:: queue_size() @@ -413,50 +279,6 @@ Reference :returns: The amount of task packages in the broker :rtype: int -.. py:function:: result_group(group_id, failures=False, wait=0, count=None, cached=False) - - Returns the results of a task group - - :param str group_id: the group identifier - :param bool failures: set this to ``True`` to include failed results - :param int wait: optional milliseconds to wait for a result or count - :param int count: block until there are this many results in the group - :param bool cached: run this against the cache backend - :returns: a list of results - :rtype: list - -.. py:function:: fetch_group(group_id, failures=True, wait=0, count=None, cached=False) - - Returns a list of tasks in a group - - :param str group_id: the group identifier - :param bool failures: set this to ``False`` to exclude failed tasks - :param int wait: optional milliseconds to wait for a task or count - :param int count: block until there are this many tasks in the group - :param bool cached: run this against the cache backend. - :returns: a list of :class:`Task` - :rtype: list - -.. py:function:: count_group(group_id, failures=False, cached=False) - - Counts the number of task results in a group. - - :param str group_id: the group identifier - :param bool failures: counts the number of failures if ``True`` - :param bool cached: run this against the cache backend. - :returns: the number of tasks or failures in a group - :rtype: int - -.. py:function:: delete_group(group_id, tasks=False, cached=False) - - Deletes a group label from the database. - - :param str group_id: the group identifier - :param bool tasks: also deletes the associated tasks if ``True`` - :param bool cached: run this against the cache backend. - :returns: the numbers of tasks affected - :rtype: int - .. py:function:: delete_cached(task_id, broker=None) Deletes a task from the cache backend @@ -576,108 +398,87 @@ Reference A proxy model of :class:`Task` with the queryset filtered on :attr:`Task.success` is ``False``. -.. py:class:: Iter(func=None, args=None, kwargs=None, cached=Conf.CACHED, sync=Conf.SYNC, broker=None) +.. py:class:: Async(func, *args, **kwargs) - An async task with iterable arguments. Serves as a convenient wrapper for :func:`async_iter` - You can pass the iterable arguments at construction or you can append individual argument tuples. + A class wrapper for the :func:`async` function. - :param func: the function to execute - :param args: an iterable of arguments. - :param kwargs: the keyword arguments - :param bool cached: run this against the cache backend - :param bool sync: execute this inline instead of asynchronous - :param broker: optional broker instance + :param object func: The task function to execute + :param tuple args: The arguments for the task function + :param dict kwargs: Keyword arguments for the task function, including async options + .. py:attribute:: id - .. py:method:: append(*args) + The task unique identifier. This will only be available after it has been :meth:`run` - Append arguments to the iter set. Returns the current set count. + .. py:attribute:: started - :param args: the arguments for a single execution - :return: the current set count - :rtype: int + Bool indicating if the task has been run with the current parameters + .. py:attribute:: func + + The task function to execute + + .. py:attribute:: args + + A tuple of arguments for the task function + + .. py:attribute:: kwargs + + Keyword arguments for the function. Can include any of the optional async keyword attributes directly or in a `q_options` dictionary. + + .. py:attribute:: broker + + Optional :class:`Broker` instance to use + + .. py:attribute:: sync + + Run this task inline instead of asynchronous. + + .. py:attribute:: save + + Overrides the global save setting. + + .. py:attribute:: hook + + Optional function to call after a result is available. Takes the result :class:`Task` as the first argument. + + .. py:attribute:: group + + Optional group identifier + + .. py:attribute:: cached + + Run the task against the cache result backend. .. py:method:: run() - Start queueing the tasks to the worker cluster. - - :return: the task result id - + Send the task to a worker cluster for execution .. py:method:: result(wait=0) - return the full list of results. + The task result. Always returns None if the task hasn't been run with the current parameters. - :param int wait: how many milliseconds to wait for a result - :return: an unsorted list of results + :param int wait: the number of milliseconds to wait for a result. -1 for indefinite .. py:method:: fetch(wait=0) - get the task result objects. + Returns the full :class:`Task` result instance. - :param int wait: how many milliseconds to wait for a result - :return: an unsorted list of task objects + :param int wait: the number of milliseconds to wait for a result. -1 for indefinite + .. py:method:: result_group(failures=False, wait=0, count=None) - .. py:method:: length() + Returns a list of results from this task's group. - get the length of the arguments list + :param bool failures: set this to ``True`` to include failed results + :param int wait: optional milliseconds to wait for a result or count. -1 for indefinite + :param int count: block until there are this many results in the group - :return int: length of the argument list + .. py:method:: fetch_group(failures=True, wait=0, count=None) + Returns a list of task results from this task's group -.. 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` - 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 str group: an optional group name. - :param bool cached: run this against the cache backend - :param bool sync: execute this inline instead of asynchronous - - - .. py:method:: append(func, *args, **kwargs) - - Append a task to the chain. Takes the same arguments as :func:`async` - - :return: the current number of tasks in the chain - :rtype: int - - - .. py:method:: run() - - Start queueing the chain to the worker cluster. - - :return: the chains group id - - - .. py:method:: result(wait=0) - - return the full list of results from the chain when it finishes. Blocks until timeout or result. - - :param int wait: how many milliseconds to wait for a result - :return: an unsorted list of results - - - .. py:method:: fetch(failures=True, wait=0) - - get the task result objects from the chain when it finishes. Blocks until timeout or result. - - :param failures: include failed tasks - :param int wait: how many milliseconds to wait for a result - :return: an unsorted list of task objects - - .. py:method:: current() - - get the index of the currently executing chain element - - :return int: current chain index - - .. py:method:: length() - - get the length of the chain - - :return int: length of the chain \ No newline at end of file + :param bool failures: set this to ``False`` to exclude failed tasks + :param int wait: optional milliseconds to wait for a task or count. -1 for indefinite + :param int count: block until there are this many tasks in the group \ No newline at end of file