mirror of
https://github.com/django-q2/django-q2.git
synced 2026-10-04 03:48:13 +08:00
+28
-17
@@ -1,21 +1,30 @@
|
|||||||
from datetime import timedelta
|
from datetime import timedelta
|
||||||
from time import sleep
|
from time import sleep
|
||||||
|
|
||||||
from django.utils import timezone
|
from django.utils import timezone
|
||||||
from django.db.models import Q
|
|
||||||
from django_q.brokers import Broker
|
from django_q.brokers import Broker
|
||||||
from django_q.models import OrmQ
|
from django_q.models import OrmQ
|
||||||
from django_q.conf import Conf
|
from django_q.conf import Conf
|
||||||
|
|
||||||
|
|
||||||
|
def _timeout():
|
||||||
|
return timezone.now() - timedelta(seconds=Conf.RETRY)
|
||||||
|
|
||||||
|
|
||||||
class ORM(Broker):
|
class ORM(Broker):
|
||||||
|
@staticmethod
|
||||||
|
def get_connection(list_key=Conf.PREFIX):
|
||||||
|
return OrmQ.objects.using(Conf.ORM)
|
||||||
|
|
||||||
def queue_size(self):
|
def queue_size(self):
|
||||||
return OrmQ.objects.using(Conf.ORM) \
|
return self.connection.filter(key=self.list_key, lock__lte=_timeout()).count()
|
||||||
.filter(Q(key=self.list_key, lock__isnull=True) |
|
|
||||||
Q(key=self.list_key, lock__lte=timezone.now() - timedelta(seconds=Conf.RETRY))) \
|
def lock_size(self):
|
||||||
.count()
|
return self.connection.filter(key=self.list_key, lock__gte=_timeout()).count()
|
||||||
|
|
||||||
def purge_queue(self):
|
def purge_queue(self):
|
||||||
return OrmQ.objects.using(Conf.ORM).filter(key=self.list_key).delete()
|
return self.connection.filter(key=self.list_key).delete()
|
||||||
|
|
||||||
def ping(self):
|
def ping(self):
|
||||||
return True
|
return True
|
||||||
@@ -27,25 +36,27 @@ class ORM(Broker):
|
|||||||
self.delete(task_id)
|
self.delete(task_id)
|
||||||
|
|
||||||
def enqueue(self, task):
|
def enqueue(self, task):
|
||||||
package = OrmQ.objects.using(Conf.ORM).create(key=self.list_key, payload=task)
|
package = self.connection.create(key=self.list_key, payload=task, lock=_timeout())
|
||||||
return package.pk
|
return package.pk
|
||||||
|
|
||||||
def dequeue(self):
|
def dequeue(self):
|
||||||
tasks = OrmQ.objects.using(Conf.ORM).filter(
|
tasks = self.connection.filter(key=self.list_key, lock__lt=_timeout())[0:Conf.BULK]
|
||||||
Q(key=self.list_key, lock__isnull=True) |
|
if tasks:
|
||||||
Q(key=self.list_key, lock__lte=timezone.now() - timedelta(seconds=Conf.RETRY)))[:Conf.BULK]
|
task_list = []
|
||||||
if tasks:
|
lock = timezone.now()
|
||||||
# lock them
|
for task in tasks:
|
||||||
OrmQ.objects.using(Conf.ORM).filter(pk__in=tasks).update(lock=timezone.now())
|
task.lock = lock
|
||||||
return [(t.pk, t.payload) for t in tasks]
|
task.save(update_fields=['lock'])
|
||||||
# empty queue, spare the cpu
|
task_list.append((task.pk, task.payload))
|
||||||
sleep(0.2)
|
return task_list
|
||||||
|
# empty queue, spare the cpu
|
||||||
|
sleep(0.2)
|
||||||
|
|
||||||
def delete_queue(self):
|
def delete_queue(self):
|
||||||
return self.purge_queue()
|
return self.purge_queue()
|
||||||
|
|
||||||
def delete(self, task_id):
|
def delete(self, task_id):
|
||||||
return OrmQ.objects.using(Conf.ORM).filter(pk=task_id).delete()
|
self.connection.filter(pk=task_id).delete()
|
||||||
|
|
||||||
def acknowledge(self, task_id):
|
def acknowledge(self, task_id):
|
||||||
return self.delete(task_id)
|
return self.delete(task_id)
|
||||||
|
|||||||
+4
-1
@@ -82,9 +82,12 @@ def monitor(run_once=False, broker=None):
|
|||||||
i += 1
|
i += 1
|
||||||
# bottom bar
|
# bottom bar
|
||||||
i += 1
|
i += 1
|
||||||
|
queue_size = broker.queue_size()
|
||||||
|
if Conf.ORM:
|
||||||
|
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)))
|
||||||
print(term.move(i, 3 * col_width) + term.white_on_cyan(term.center(broker.queue_size(), width=col_width)))
|
print(term.move(i, 3 * col_width) + term.white_on_cyan(term.center(queue_size, width=col_width)))
|
||||||
print(term.move(i, 4 * col_width) + term.black_on_cyan(term.center(_('Success'), width=col_width)))
|
print(term.move(i, 4 * col_width) + term.black_on_cyan(term.center(_('Success'), width=col_width)))
|
||||||
print(term.move(i, 5 * col_width) + term.white_on_cyan(
|
print(term.move(i, 5 * col_width) + term.white_on_cyan(
|
||||||
term.center(models.Success.objects.count(), width=col_width)))
|
term.center(models.Success.objects.count(), width=col_width)))
|
||||||
|
|||||||
@@ -268,9 +268,12 @@ def test_orm():
|
|||||||
broker.enqueue('test')
|
broker.enqueue('test')
|
||||||
Conf.BULK = 5
|
Conf.BULK = 5
|
||||||
tasks = broker.dequeue()
|
tasks = broker.dequeue()
|
||||||
|
assert broker.lock_size() == Conf.BULK
|
||||||
for task in tasks:
|
for task in tasks:
|
||||||
assert task is not None
|
assert task is not None
|
||||||
broker.acknowledge(task[0])
|
broker.acknowledge(task[0])
|
||||||
|
# test lock size
|
||||||
|
assert broker.lock_size() == 0
|
||||||
# test duplicate acknowledge
|
# test duplicate acknowledge
|
||||||
broker.acknowledge(task[0])
|
broker.acknowledge(task[0])
|
||||||
# delete queue
|
# delete queue
|
||||||
|
|||||||
@@ -4,6 +4,7 @@ from django_q import async
|
|||||||
from django_q.cluster import Cluster
|
from django_q.cluster import Cluster
|
||||||
from django_q.monitor import monitor, info
|
from django_q.monitor import monitor, info
|
||||||
from django_q.status import Stat
|
from django_q.status import Stat
|
||||||
|
from django_q.conf import Conf
|
||||||
|
|
||||||
|
|
||||||
@pytest.mark.django_db
|
@pytest.mark.django_db
|
||||||
@@ -22,6 +23,10 @@ def test_monitor():
|
|||||||
assert stat.empty_queues() is True
|
assert stat.empty_queues() is True
|
||||||
break
|
break
|
||||||
assert found_c is True
|
assert found_c is True
|
||||||
|
# test lock size for orm broker
|
||||||
|
Conf.ORM = 'default'
|
||||||
|
monitor(run_once=True)
|
||||||
|
Conf.ORM = None
|
||||||
|
|
||||||
|
|
||||||
@pytest.mark.django_db
|
@pytest.mark.django_db
|
||||||
|
|||||||
Reference in New Issue
Block a user