Merge pull request #68 from Koed00/dev

Adds a wait option for results.
This commit is contained in:
Ilan Steemers
2015-09-18 19:31:35 +02:00
9 changed files with 49 additions and 12 deletions
+4 -1
View File
@@ -127,7 +127,10 @@ Use `async` from your code to quickly offload tasks:
task_result = result(task_id) task_result = result(task_id)
# result returns None if the task has not been executed yet # result returns None if the task has not been executed yet
# so in most cases you will want to use a hook: # you can wait for it
task_result = result(task_id, 200)
# but in most cases you will want to use a hook:
async('math.modf', 2.5, hook='hooks.print_result') async('math.modf', 2.5, hook='hooks.print_result')
+3
View File
@@ -27,6 +27,9 @@ class Sqs(Broker):
def queue_size(self): def queue_size(self):
return int(self.queue.attributes['ApproximateNumberOfMessages']) return int(self.queue.attributes['ApproximateNumberOfMessages'])
def lock_size(self):
return int(self.queue.attributes['ApproximateNumberOfMessagesNotVisible'])
def delete(self, task_id): def delete(self, task_id):
message = self.sqs.Message(self.queue.url, task_id) message = self.sqs.Message(self.queue.url, task_id)
message.delete() message.delete()
+1 -2
View File
@@ -83,7 +83,7 @@ def monitor(run_once=False, broker=None):
# bottom bar # bottom bar
i += 1 i += 1
queue_size = broker.queue_size() queue_size = broker.queue_size()
if Conf.ORM: if hasattr(broker, 'lock_size'):
queue_size = '{}({})'.format(queue_size, broker.lock_size()) queue_size = '{}({})'.format(queue_size, broker.lock_size())
print(term.move(i, 0) + term.white_on_cyan(term.center(broker.info(), width=col_width * 2))) print(term.move(i, 0) + term.white_on_cyan(term.center(broker.info(), width=col_width * 2)))
print(term.move(i, 2 * col_width) + term.black_on_cyan(term.center(_('Queued'), width=col_width))) print(term.move(i, 2 * col_width) + term.black_on_cyan(term.center(_('Queued'), width=col_width)))
@@ -183,4 +183,3 @@ def info(broker=None):
term.white('{0:.4f}'.format(exec_time)) term.white('{0:.4f}'.format(exec_time))
) )
return True return True
+23 -4
View File
@@ -5,6 +5,7 @@ from multiprocessing import Queue, Value
from django.utils import timezone from django.utils import timezone
# local # local
import time
import signing import signing
import cluster import cluster
from django_q.conf import Conf, logger from django_q.conf import Conf, logger
@@ -82,16 +83,25 @@ def schedule(func, *args, **kwargs):
) )
def result(task_id): def result(task_id, wait=0):
""" """
Return the result of the named task. Return the result of the named task.
:type task_id: str or uuid :type task_id: str or uuid
:param task_id: the task name or uuid :param task_id: the task name or uuid
:type wait: int
:param wait: number of milliseconds to wait for a result
:return: the result object of this task :return: the result object of this task
:rtype: object :rtype: object
""" """
return Task.get_result(task_id) start = time.time()
while True:
r = Task.get_result(task_id)
if r:
return r
if (time.time() - start) * 1000 >= wait:
break
time.sleep(0.01)
def result_group(group_id, failures=False): def result_group(group_id, failures=False):
@@ -105,16 +115,25 @@ def result_group(group_id, failures=False):
return Task.get_result_group(group_id, failures) return Task.get_result_group(group_id, failures)
def fetch(task_id): def fetch(task_id, wait=0):
""" """
Return the processed task. Return the processed task.
:param task_id: the task name or uuid :param task_id: the task name or uuid
:type task_id: str or uuid :type task_id: str or uuid
:param wait: the number of milliseconds to wait for a result
:type wait: int
:return: the full task object :return: the full task object
:rtype: Task :rtype: Task
""" """
return Task.get_task(task_id) start = time.time()
while True:
t = Task.get_task(task_id)
if t:
return t
if (time.time() - start) * 1000 >= wait:
break
time.sleep(0.01)
def fetch_group(group_id, failures=True): def fetch_group(group_id, failures=True):
+1
View File
@@ -216,6 +216,7 @@ def test_sqs():
broker.acknowledge(task[0]) broker.acknowledge(task[0])
# duplicate acknowledge # duplicate acknowledge
broker.acknowledge(task[0]) broker.acknowledge(task[0])
assert broker.lock_size() == 0
# delete queue # delete queue
broker.enqueue('test') broker.enqueue('test')
broker.purge_queue() broker.purge_queue()
+2
View File
@@ -230,6 +230,8 @@ def test_async(broker, admin_user):
assert result_j.group_delete(tasks=True) is None assert result_j.group_delete(tasks=True) is None
# task k should not have been saved # task k should not have been saved
assert fetch(k) is None assert fetch(k) is None
assert fetch(k, 100) is None
assert result(k, 100) is None
broker.delete_queue() broker.delete_queue()
+5
View File
@@ -129,6 +129,11 @@ You can override this class if you want to contribute and support your own broke
Returns the amount of messages in the brokers queue. Returns the amount of messages in the brokers queue.
.. py:method:: lock_size()
Optional method that returns the number of messages currently awaiting acknowledgement.
Only implemented on brokers that support it.
.. py:method:: ping() .. py:method:: ping()
Returns True if the broker can be reached. Returns True if the broker can be reached.
+9 -4
View File
@@ -25,7 +25,10 @@ Use :func:`async` from your code to quickly offload tasks to the :class:`Cluster
task_result = result(task_id) task_result = result(task_id)
# result returns None if the task has not been executed yet # result returns None if the task has not been executed yet
# so in most cases you will want to use a hook: # you can wait for it
task_result = result(task_id, 200)
# but in most cases you will want to use a hook:
async('math.modf', 2.5, hook='hooks.print_result') async('math.modf', 2.5, hook='hooks.print_result')
@@ -208,19 +211,21 @@ Reference
:returns: The uuid of the task :returns: The uuid of the task
:rtype: str :rtype: str
.. py:function:: result(task_id) .. py:function:: result(task_id, wait=0)
Gets the result of a previously executed task Gets the result of a previously executed task
:param str task_id: the uuid or name of the task :param str task_id: the uuid or name of the task
:param int wait: optional milliseconds to wait for a result
:returns: The result of the executed task :returns: The result of the executed task
.. py:function:: fetch(task_id) .. py:function:: fetch(task_id, wait=0)
Returns a previously executed task Returns a previously executed task
:param str name: the uuid or name of the task :param str name: the uuid or name of the task
:returns: The task if any :param int wait: optional milliseconds to wait for a result
:returns: A task object
:rtype: Task :rtype: Task
.. versionchanged:: 0.2.0 .. versionchanged:: 0.2.0
+1 -1
View File
@@ -7,7 +7,7 @@
arrow==0.6.0 arrow==0.6.0
blessed==1.9.5 blessed==1.9.5
boto3==1.1.3 boto3==1.1.3
botocore==1.2.2 # via boto3 botocore==1.2.4 # via boto3
django-picklefield==0.3.2 django-picklefield==0.3.2
django-redis==4.2.0 django-redis==4.2.0
docutils==0.12 # via botocore docutils==0.12 # via botocore