mirror of
https://github.com/django-q2/django-q2.git
synced 2026-09-15 21:47:53 +08:00
Django ORM broker
This commit is contained in:
@@ -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)
|
||||
|
||||
@@ -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
|
||||
|
||||
62
django_q/brokers/orm.py
Normal file
62
django_q/brokers/orm.py
Normal file
@@ -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)
|
||||
@@ -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
|
||||
|
||||
@@ -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:
|
||||
|
||||
27
django_q/migrations/0007_ormq.py
Normal file
27
django_q/migrations/0007_ormq.py
Normal file
@@ -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',
|
||||
},
|
||||
),
|
||||
]
|
||||
@@ -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':
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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()
|
||||
|
||||
|
||||
@@ -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()
|
||||
|
||||
|
||||
Reference in New Issue
Block a user