From 2bffcf554b962d5e38cc689afac6f2ab690cdebf Mon Sep 17 00:00:00 2001 From: Ilan Steemers Date: Mon, 14 Sep 2015 13:24:43 +0200 Subject: [PATCH] Django ORM broker --- django_q/admin.py | 27 ++++++++++---- django_q/brokers/__init__.py | 3 ++ django_q/brokers/orm.py | 62 ++++++++++++++++++++++++++++++++ django_q/cluster.py | 2 +- django_q/conf.py | 6 ++++ django_q/migrations/0007_ormq.py | 27 ++++++++++++++ django_q/models.py | 11 ++++++ django_q/tests/test_brokers.py | 55 ++++++++++++++++++++++++++++ django_q/tests/test_cluster.py | 1 + django_q/tests/test_scheduler.py | 1 + 10 files changed, 187 insertions(+), 8 deletions(-) create mode 100644 django_q/brokers/orm.py create mode 100644 django_q/migrations/0007_ormq.py diff --git a/django_q/admin.py b/django_q/admin.py index f03e90f..e8a725d 100644 --- a/django_q/admin.py +++ b/django_q/admin.py @@ -3,11 +3,11 @@ from django.contrib import admin from django.utils.translation import ugettext_lazy as _ from .tasks import async -from .models import Success, Failure, Schedule +from .models import Success, Failure, Schedule, OrmQ +from .conf import Conf class TaskAdmin(admin.ModelAdmin): - """model admin for success tasks.""" list_display = ( @@ -34,8 +34,8 @@ class TaskAdmin(admin.ModelAdmin): def get_readonly_fields(self, request, obj=None): """Set all fields readonly.""" - return list(self.readonly_fields) +\ - [field.name for field in obj._meta.fields] + return list(self.readonly_fields) + \ + [field.name for field in obj._meta.fields] def retry_failed(FailAdmin, request, queryset): @@ -49,7 +49,6 @@ retry_failed.short_description = _("Resubmit selected tasks to queue") class FailAdmin(admin.ModelAdmin): - """model admin for failed tasks.""" list_display = ( @@ -72,11 +71,10 @@ class FailAdmin(admin.ModelAdmin): def get_readonly_fields(self, request, obj=None): """Set all fields readonly.""" return list(self.readonly_fields) + \ - [field.name for field in obj._meta.fields] + [field.name for field in obj._meta.fields] class ScheduleAdmin(admin.ModelAdmin): - """ model admin for schedules """ list_display = ( @@ -95,6 +93,21 @@ class ScheduleAdmin(admin.ModelAdmin): list_display_links = ('id', 'name') +class QueueAdmin(admin.ModelAdmin): + """ queue admin for ORM broker """ + list_display = ( + 'id', + 'key', + 'lock' + ) + + def has_add_permission(self, request, obj=None): + """Don't allow adds.""" + return False + admin.site.register(Schedule, ScheduleAdmin) admin.site.register(Success, TaskAdmin) admin.site.register(Failure, FailAdmin) + +if Conf.ORM: + admin.site.register(OrmQ, QueueAdmin) diff --git a/django_q/brokers/__init__.py b/django_q/brokers/__init__.py index fed4590..0b5b84f 100644 --- a/django_q/brokers/__init__.py +++ b/django_q/brokers/__init__.py @@ -160,6 +160,9 @@ def get_broker(list_key=Conf.PREFIX): elif Conf.SQS: from brokers import aws_sqs return aws_sqs.Sqs(list_key=list_key) + elif Conf.ORM: + from brokers import orm + return orm.ORM(list_key=list_key) # default to redis else: from brokers import redis_broker diff --git a/django_q/brokers/orm.py b/django_q/brokers/orm.py new file mode 100644 index 0000000..461f894 --- /dev/null +++ b/django_q/brokers/orm.py @@ -0,0 +1,62 @@ +from datetime import timedelta +from time import sleep +from django.utils import timezone +from django.db.models import Q +from django_q.brokers import Broker +from django_q.models import OrmQ +from django_q.conf import Conf + + +class ORM(Broker): + def queue_size(self): + return 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))) \ + .count() + + def purge_queue(self): + return OrmQ.objects.using(Conf.ORM).filter(key=self.list_key).delete() + + def ping(self): + return True + + def info(self): + return 'ORM {}'.format(Conf.ORM) + + def fail(self, task_id): + self.delete(task_id) + + def enqueue(self, task): + package = OrmQ.objects.using(Conf.ORM).create(key=self.list_key, payload=task) + 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 + # empty queue, spare the cpu + sleep(0.2) + + def delete_queue(self): + return self.purge_queue() + + def delete(self, task_id): + return OrmQ.objects.using(Conf.ORM).filter(pk=task_id).delete() + + def acknowledge(self, task_id): + return self.delete(task_id) diff --git a/django_q/cluster.py b/django_q/cluster.py index 13142a6..8e99fab 100644 --- a/django_q/cluster.py +++ b/django_q/cluster.py @@ -238,7 +238,7 @@ class Sentinel(object): self.reincarnate(self.pusher) # Call scheduler once a minute (or so) counter += cycle - if counter == 30: + if counter == 30 and Conf.SCHEDULER: counter = 0 scheduler(broker=self.broker) # Save current status diff --git a/django_q/conf.py b/django_q/conf.py index 3984a28..b312c4f 100644 --- a/django_q/conf.py +++ b/django_q/conf.py @@ -47,6 +47,9 @@ class Conf(object): # SQS broker SQS = conf.get('sqs', None) + # ORM broker + ORM = conf.get('orm', None) + # Name of the cluster or site. For when you run multiple sites on one redis server PREFIX = conf.get('name', 'default') @@ -57,6 +60,9 @@ class Conf(object): # Failures are always saved SAVE_LIMIT = conf.get('save_limit', 250) + # Disable the scheduler + SCHEDULER = conf.get('scheduler', True) + # Number of workers in the pool. Default is cpu count if implemented, otherwise 4. WORKERS = conf.get('workers', False) if not WORKERS: diff --git a/django_q/migrations/0007_ormq.py b/django_q/migrations/0007_ormq.py new file mode 100644 index 0000000..1c42dc2 --- /dev/null +++ b/django_q/migrations/0007_ormq.py @@ -0,0 +1,27 @@ +# -*- coding: utf-8 -*- +from __future__ import unicode_literals + +from django.db import models, migrations + + +class Migration(migrations.Migration): + + dependencies = [ + ('django_q', '0006_auto_20150805_1817'), + ] + + operations = [ + migrations.CreateModel( + name='OrmQ', + fields=[ + ('id', models.AutoField(primary_key=True, auto_created=True, verbose_name='ID', serialize=False)), + ('key', models.CharField(max_length=100)), + ('payload', models.TextField()), + ('lock', models.DateTimeField(null=True)), + ], + options={ + 'verbose_name_plural': 'Queued tasks', + 'verbose_name': 'Queued task', + }, + ), + ] diff --git a/django_q/models.py b/django_q/models.py index d42d4af..de05408 100644 --- a/django_q/models.py +++ b/django_q/models.py @@ -188,6 +188,17 @@ class Schedule(models.Model): ordering = ['next_run'] +class OrmQ(models.Model): + key = models.CharField(max_length=100) + payload = models.TextField() + lock = models.DateTimeField(null=True) + + class Meta: + app_label = 'django_q' + verbose_name = _('Queued task') + verbose_name_plural = _('Queued tasks') + + # Backwards compatibility for Django 1.7 def decode_results(values): if get_version().split('.')[1] == '7': diff --git a/django_q/tests/test_brokers.py b/django_q/tests/test_brokers.py index 96c07f8..350c2a0 100644 --- a/django_q/tests/test_brokers.py +++ b/django_q/tests/test_brokers.py @@ -215,3 +215,58 @@ def test_sqs(): Conf.SQS = None Conf.BULK = 1 Conf.DJANGO_REDIS = 'default' + +@pytest.mark.django_db +def test_orm(): + Conf.ORM = 'default' + # check broker + broker = get_broker(list_key='orm_test') + assert broker.ping() is True + assert broker.info() is not None + # clear before we start + broker.delete_queue() + # enqueue + broker.enqueue('test') + assert broker.queue_size() == 1 + # dequeue + task = broker.dequeue() + assert task[1] == 'test' + broker.acknowledge(task[0]) + assert broker.queue_size() == 0 + # Retry test + Conf.RETRY = 1 + broker.enqueue('test') + assert broker.queue_size() == 1 + broker.dequeue() + assert broker.queue_size() == 0 + sleep(1.5) + assert broker.queue_size() == 1 + task = broker.dequeue() + assert broker.queue_size() == 0 + broker.acknowledge(task[0]) + sleep(1.5) + assert broker.queue_size() == 0 + # delete job + task_id = broker.enqueue('test') + broker.delete(task_id) + assert broker.dequeue() is None + # fail + task_id = broker.enqueue('test') + broker.fail(task_id) + # bulk test + for i in range(5): + broker.enqueue('test') + Conf.BULK = 5 + for i in range(5): + task = broker.dequeue() + assert task is not None + broker.acknowledge(task[0]) + # test duplicate acknowledge + broker.acknowledge(task[0]) + # delete queue + broker.enqueue('test') + broker.enqueue('test') + broker.delete_queue() + assert broker.queue_size() == 0 + # back to django-redis + Conf.ORM = None diff --git a/django_q/tests/test_cluster.py b/django_q/tests/test_cluster.py index dcd5394..ceabb72 100644 --- a/django_q/tests/test_cluster.py +++ b/django_q/tests/test_cluster.py @@ -32,6 +32,7 @@ def broker(): Conf.DISQUE_NODES = None Conf.IRON_MQ = None Conf.SQS = None + Conf.ORM = None Conf.DJANGO_REDIS = 'default' return get_broker() diff --git a/django_q/tests/test_scheduler.py b/django_q/tests/test_scheduler.py index 4cf7fd3..a8e5abc 100644 --- a/django_q/tests/test_scheduler.py +++ b/django_q/tests/test_scheduler.py @@ -17,6 +17,7 @@ def broker(): Conf.DISQUE_NODES = None Conf.IRON_MQ = None Conf.SQS = None + Conf.ORM = None Conf.DJANGO_REDIS = 'default' return get_broker()