diff --git a/django_q/brokers/__init__.py b/django_q/brokers/__init__.py index 0b5b84f..a8380ee 100644 --- a/django_q/brokers/__init__.py +++ b/django_q/brokers/__init__.py @@ -7,7 +7,6 @@ class Broker(object): self.connection = self.get_connection(list_key) self.list_key = list_key self.cache = self.get_cache() - self.task_cache = [] def enqueue(self, task): """ diff --git a/django_q/brokers/aws_sqs.py b/django_q/brokers/aws_sqs.py index db25215..2334c20 100644 --- a/django_q/brokers/aws_sqs.py +++ b/django_q/brokers/aws_sqs.py @@ -17,17 +17,9 @@ class Sqs(Broker): # sqs supports max 10 messages in bulk if Conf.BULK > 10: Conf.BULK = 10 - t = None - if len(self.task_cache) > 0: - t = self.task_cache.pop() - else: - tasks = self.queue.receive_messages(MaxNumberOfMessages=Conf.BULK, VisibilityTimeout=Conf.RETRY) - if tasks: - t = tasks.pop() - if tasks: - self.task_cache = tasks - if t: - return t.receipt_handle, t.body + tasks = self.queue.receive_messages(MaxNumberOfMessages=Conf.BULK, VisibilityTimeout=Conf.RETRY) + if tasks: + return [(t.receipt_handle, t.body) for t in tasks] def acknowledge(self, task_id): return self.delete(task_id) diff --git a/django_q/brokers/disque.py b/django_q/brokers/disque.py index fb0edd2..3027c1a 100644 --- a/django_q/brokers/disque.py +++ b/django_q/brokers/disque.py @@ -11,18 +11,10 @@ class Disque(Broker): 'ADDJOB {} {} 500 RETRY {}'.format(self.list_key, task, retry)).decode() def dequeue(self): - t = None - if len(self.task_cache) > 0: - t = self.task_cache.pop() - else: tasks = self.connection.execute_command( 'GETJOB COUNT {} TIMEOUT 1000 FROM {}'.format(Conf.BULK, self.list_key)) if tasks: - t = tasks.pop() - if tasks: - self.task_cache = tasks - if t: - return t[1].decode(), t[2].decode() + return [(t[1].decode(), t[2].decode()) for t in tasks] def queue_size(self): return self.connection.execute_command('QLEN {}'.format(self.list_key)) diff --git a/django_q/brokers/ironmq.py b/django_q/brokers/ironmq.py index 08bb860..f0b0725 100644 --- a/django_q/brokers/ironmq.py +++ b/django_q/brokers/ironmq.py @@ -9,18 +9,10 @@ class IronMQBroker(Broker): return self.connection.post(task)['ids'][0] def dequeue(self): - t = None - if len(self.task_cache) > 0: - t = self.task_cache.pop() - else: timeout = Conf.RETRY or None tasks = self.connection.get(timeout=timeout, wait=1, max=Conf.BULK)['messages'] if tasks: - t = tasks.pop() - if tasks: - self.task_cache = tasks - if t: - return t['id'], t['body'] + return [(t['id'], t['body']) for t in tasks] def ping(self): return self.connection.name == self.list_key diff --git a/django_q/brokers/orm.py b/django_q/brokers/orm.py index 461f894..b56c560 100644 --- a/django_q/brokers/orm.py +++ b/django_q/brokers/orm.py @@ -31,24 +31,13 @@ class ORM(Broker): return package.pk def dequeue(self): - if len(self.task_cache) > 0: - t = self.task_cache.pop() - return t.pk, t.payload - else: - # Get new and timed out tasks tasks = OrmQ.objects.using(Conf.ORM).filter( Q(key=self.list_key, lock__isnull=True) | Q(key=self.list_key, lock__lte=timezone.now() - timedelta(seconds=Conf.RETRY)))[:Conf.BULK] if tasks: # lock them OrmQ.objects.using(Conf.ORM).filter(pk__in=tasks).update(lock=timezone.now()) - tasks = [t for t in tasks] - # pop one task - t = tasks.pop() - if tasks: - # add remainder to cache - self.task_cache = [t for t in tasks] - return t.pk, t.payload + return [(t.pk, t.payload) for t in tasks] # empty queue, spare the cpu sleep(0.2) diff --git a/django_q/brokers/redis_broker.py b/django_q/brokers/redis_broker.py index 68289aa..6a8cf9a 100644 --- a/django_q/brokers/redis_broker.py +++ b/django_q/brokers/redis_broker.py @@ -19,7 +19,7 @@ class Redis(Broker): def dequeue(self): task = self.connection.blpop(self.list_key, 1) if task: - return None, task[1] + return [(None, task[1])] def queue_size(self): return self.connection.llen(self.list_key) diff --git a/django_q/cluster.py b/django_q/cluster.py index 8e99fab..e596267 100644 --- a/django_q/cluster.py +++ b/django_q/cluster.py @@ -298,23 +298,24 @@ def pusher(task_queue, event, broker=None): logger.info(_('{} pushing tasks at {}').format(current_process().name, current_process().pid)) while True: try: - task = broker.dequeue() + task_set = broker.dequeue() except Exception as e: logger.error(e) # broker probably crashed. Let the sentinel handle it. sleep(10) break - if task: - ack_id = task[0] - # unpack the task - try: - task = signing.SignedPackage.loads(task[1]) - except (TypeError, signing.BadSignature) as e: - logger.error(e) - broker.fail(ack_id) - continue - task['ack_id'] = ack_id - task_queue.put(task) + if task_set: + for task in task_set: + ack_id = task[0] + # unpack the task + try: + task = signing.SignedPackage.loads(task[1]) + except (TypeError, signing.BadSignature) as e: + logger.error(e) + broker.fail(ack_id) + continue + task['ack_id'] = ack_id + task_queue.put(task) logger.debug(_('queueing from {}').format(broker.list_key)) if event.is_set(): break diff --git a/django_q/tests/test_brokers.py b/django_q/tests/test_brokers.py index 350c2a0..68d8534 100644 --- a/django_q/tests/test_brokers.py +++ b/django_q/tests/test_brokers.py @@ -49,7 +49,7 @@ def test_disque(): broker.enqueue('test') assert broker.queue_size() == 1 # dequeue - task = broker.dequeue() + task = broker.dequeue()[0] assert task[1] == 'test' broker.acknowledge(task[0]) assert broker.queue_size() == 0 @@ -61,7 +61,7 @@ def test_disque(): assert broker.queue_size() == 0 sleep(1.5) assert broker.queue_size() == 1 - task = broker.dequeue() + task = broker.dequeue()[0] assert broker.queue_size() == 0 broker.acknowledge(task[0]) sleep(1.5) @@ -78,8 +78,8 @@ def test_disque(): broker.enqueue('test') Conf.BULK = 5 Conf.DISQUE_FASTACK = True - for i in range(5): - task = broker.dequeue() + tasks = broker.dequeue() + for task in tasks: assert task is not None broker.acknowledge(task[0]) # test duplicate acknowledge @@ -117,7 +117,7 @@ def test_ironmq(): # enqueue broker.enqueue('test') # dequeue - task = broker.dequeue() + task = broker.dequeue()[0] assert task[1] == 'test' broker.acknowledge(task[0]) assert broker.dequeue() is None @@ -126,7 +126,7 @@ def test_ironmq(): broker.enqueue('test') assert broker.dequeue() is not None sleep(1.5) - task = broker.dequeue() + task = broker.dequeue()[0] assert len(task) > 0 broker.acknowledge(task[0]) sleep(1.5) @@ -141,8 +141,8 @@ def test_ironmq(): for i in range(5): broker.enqueue('test') Conf.BULK = 5 - for i in range(5): - task = broker.dequeue() + tasks = broker.dequeue() + for task in tasks: assert task is not None broker.acknowledge(task[0]) # duplicate acknowledge @@ -174,7 +174,7 @@ def test_sqs(): # enqueue broker.enqueue('test') # dequeue - task = broker.dequeue() + task = broker.dequeue()[0] assert task[1] == 'test' broker.acknowledge(task[0]) assert broker.dequeue() is None @@ -183,26 +183,26 @@ def test_sqs(): broker.enqueue('test') assert broker.dequeue() is not None sleep(1.5) - task = broker.dequeue() + task = broker.dequeue()[0] assert len(task) > 0 broker.acknowledge(task[0]) sleep(1.5) # delete job broker.enqueue('test') - task_id = broker.dequeue()[0] + task_id = broker.dequeue()[0][0] broker.delete(task_id) assert broker.dequeue() is None # fail broker.enqueue('test') while task is None: - task = broker.dequeue() + task = broker.dequeue()[0] broker.fail(task[0]) # bulk test for i in range(10): broker.enqueue('test') Conf.BULK = 12 - for i in range(10): - task = broker.dequeue() + tasks = broker.dequeue() + for task in tasks: assert task is not None broker.acknowledge(task[0]) # duplicate acknowledge @@ -216,6 +216,7 @@ def test_sqs(): Conf.BULK = 1 Conf.DJANGO_REDIS = 'default' + @pytest.mark.django_db def test_orm(): Conf.ORM = 'default' @@ -229,7 +230,7 @@ def test_orm(): broker.enqueue('test') assert broker.queue_size() == 1 # dequeue - task = broker.dequeue() + task = broker.dequeue()[0] assert task[1] == 'test' broker.acknowledge(task[0]) assert broker.queue_size() == 0 @@ -241,7 +242,7 @@ def test_orm(): assert broker.queue_size() == 0 sleep(1.5) assert broker.queue_size() == 1 - task = broker.dequeue() + task = broker.dequeue()[0] assert broker.queue_size() == 0 broker.acknowledge(task[0]) sleep(1.5) @@ -257,8 +258,8 @@ def test_orm(): for i in range(5): broker.enqueue('test') Conf.BULK = 5 - for i in range(5): - task = broker.dequeue() + tasks = broker.dequeue() + for task in tasks: assert task is not None broker.acknowledge(task[0]) # test duplicate acknowledge