mirror of
https://github.com/django-q2/django-q2.git
synced 2026-10-04 01:08:13 +08:00
IronMQ broker
This commit is contained in:
@@ -9,6 +9,6 @@ from .models import Task, Schedule, Success, Failure
|
|||||||
from .cluster import Cluster
|
from .cluster import Cluster
|
||||||
from .status import Stat
|
from .status import Stat
|
||||||
|
|
||||||
VERSION = (0, 6, 2)
|
VERSION = (0, 6, 3)
|
||||||
|
|
||||||
default_app_config = 'django_q.apps.DjangoQConfig'
|
default_app_config = 'django_q.apps.DjangoQConfig'
|
||||||
|
|||||||
@@ -153,6 +153,9 @@ def get_broker(list_key=Conf.PREFIX):
|
|||||||
if Conf.DISQUE_NODES:
|
if Conf.DISQUE_NODES:
|
||||||
from brokers import disque
|
from brokers import disque
|
||||||
return disque.Disque(list_key=list_key)
|
return disque.Disque(list_key=list_key)
|
||||||
|
elif Conf.IRON_MQ:
|
||||||
|
from brokers import ironmq
|
||||||
|
return ironmq.IronMQBroker(list_key=list_key)
|
||||||
# default to redis
|
# default to redis
|
||||||
else:
|
else:
|
||||||
from brokers import redis_broker
|
from brokers import redis_broker
|
||||||
|
|||||||
@@ -0,0 +1,44 @@
|
|||||||
|
from django_q.conf import Conf
|
||||||
|
from django_q.brokers import Broker
|
||||||
|
from iron_mq import IronMQ
|
||||||
|
|
||||||
|
|
||||||
|
class IronMQBroker(Broker):
|
||||||
|
|
||||||
|
def enqueue(self, task):
|
||||||
|
return self.connection.post(task)['ids'][0]
|
||||||
|
|
||||||
|
def dequeue(self):
|
||||||
|
timeout = Conf.RETRY or None
|
||||||
|
task = self.connection.get(timeout=timeout, wait=1)['messages']
|
||||||
|
if task:
|
||||||
|
return task[0]['id'], task[0]['body']
|
||||||
|
|
||||||
|
def ping(self):
|
||||||
|
return self.connection.name == self.list_key
|
||||||
|
|
||||||
|
def info(self):
|
||||||
|
return 'IronMQ'
|
||||||
|
|
||||||
|
def queue_size(self):
|
||||||
|
return self.connection.size()
|
||||||
|
|
||||||
|
def delete_queue(self):
|
||||||
|
return self.connection.delete_queue()['msg']
|
||||||
|
|
||||||
|
def purge_queue(self):
|
||||||
|
return self.connection.clear()
|
||||||
|
|
||||||
|
def delete(self, task_id):
|
||||||
|
return self.connection.delete(task_id)['msg']
|
||||||
|
|
||||||
|
def fail(self, task_id):
|
||||||
|
self.delete(task_id)
|
||||||
|
|
||||||
|
def acknowledge(self, task_id):
|
||||||
|
return self.delete(task_id)
|
||||||
|
|
||||||
|
@staticmethod
|
||||||
|
def get_connection(list_key=Conf.PREFIX):
|
||||||
|
ironmq = IronMQ(name=None, **Conf.IRON_MQ)
|
||||||
|
return ironmq.queue(queue_name=list_key)
|
||||||
+4
-1
@@ -37,6 +37,9 @@ class Conf(object):
|
|||||||
# Optional Authentication
|
# Optional Authentication
|
||||||
DISQUE_AUTH = conf.get('disque_auth', None)
|
DISQUE_AUTH = conf.get('disque_auth', None)
|
||||||
|
|
||||||
|
# IronMQ broker
|
||||||
|
IRON_MQ = conf.get('iron_mq', None)
|
||||||
|
|
||||||
# Name of the cluster or site. For when you run multiple sites on one redis server
|
# Name of the cluster or site. For when you run multiple sites on one redis server
|
||||||
PREFIX = conf.get('name', 'default')
|
PREFIX = conf.get('name', 'default')
|
||||||
|
|
||||||
@@ -62,7 +65,7 @@ class Conf(object):
|
|||||||
WORKERS = 4
|
WORKERS = 4
|
||||||
|
|
||||||
# Maximum number of tasks that each cluster can work on
|
# Maximum number of tasks that each cluster can work on
|
||||||
QUEUE_LIMIT = conf.get('queue_limit', int(WORKERS)**2)
|
QUEUE_LIMIT = conf.get('queue_limit', int(WORKERS) ** 2)
|
||||||
|
|
||||||
# Sets compression of redis packages
|
# Sets compression of redis packages
|
||||||
COMPRESSED = conf.get('compress', False)
|
COMPRESSED = conf.get('compress', False)
|
||||||
|
|||||||
@@ -74,7 +74,7 @@ def test_disque():
|
|||||||
broker.delete(task_id)
|
broker.delete(task_id)
|
||||||
assert broker.queue_size() == 0
|
assert broker.queue_size() == 0
|
||||||
# fail
|
# fail
|
||||||
task_id=broker.enqueue('test')
|
task_id = broker.enqueue('test')
|
||||||
broker.fail(task_id)
|
broker.fail(task_id)
|
||||||
# delete queue
|
# delete queue
|
||||||
broker.enqueue('test')
|
broker.enqueue('test')
|
||||||
@@ -83,3 +83,49 @@ def test_disque():
|
|||||||
assert broker.queue_size() == 0
|
assert broker.queue_size() == 0
|
||||||
# back to django-redis
|
# back to django-redis
|
||||||
Conf.DISQUE_NODES = None
|
Conf.DISQUE_NODES = None
|
||||||
|
|
||||||
|
|
||||||
|
@pytest.mark.skipif(not os.getenv('IRON_MQ_TOKEN'),
|
||||||
|
reason="requires IronMQ credentials")
|
||||||
|
def test_ironmq():
|
||||||
|
Conf.IRON_MQ = {'host': os.getenv('IRON_MQ_HOST'),
|
||||||
|
'token': os.getenv('IRON_MQ_TOKEN'),
|
||||||
|
'project_id': os.getenv('IRON_MQ_PROJECT_ID')}
|
||||||
|
# check broker
|
||||||
|
broker = get_broker(list_key='djangoQ')
|
||||||
|
assert broker.ping() is True
|
||||||
|
assert broker.info() is not None
|
||||||
|
# clear before we start
|
||||||
|
broker.purge_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
|
||||||
|
assert broker.dequeue() is not None
|
||||||
|
sleep(1.5)
|
||||||
|
task = broker.dequeue()
|
||||||
|
assert len(task) > 0
|
||||||
|
broker.acknowledge(task[0])
|
||||||
|
sleep(1.5)
|
||||||
|
# delete job
|
||||||
|
task_id = broker.enqueue('test')
|
||||||
|
broker.delete(task_id)
|
||||||
|
assert broker.queue_size() == 0
|
||||||
|
# fail
|
||||||
|
task_id = broker.enqueue('test')
|
||||||
|
broker.fail(task_id)
|
||||||
|
# delete queue
|
||||||
|
broker.enqueue('test')
|
||||||
|
broker.enqueue('test')
|
||||||
|
broker.purge_queue()
|
||||||
|
assert broker.queue_size() == 0
|
||||||
|
# back to django-redis
|
||||||
|
Conf.IRON_MQ = None
|
||||||
|
|||||||
+1
-1
@@ -72,7 +72,7 @@ author = 'Ilan Steemers'
|
|||||||
# The short X.Y version.
|
# The short X.Y version.
|
||||||
version = '0.6'
|
version = '0.6'
|
||||||
# The full version, including alpha/beta/rc tags.
|
# The full version, including alpha/beta/rc tags.
|
||||||
release = '0.6.2'
|
release = '0.6.3'
|
||||||
|
|
||||||
# The language for content autogenerated by Sphinx. Refer to documentation
|
# The language for content autogenerated by Sphinx. Refer to documentation
|
||||||
# for a list of supported languages.
|
# for a list of supported languages.
|
||||||
|
|||||||
@@ -6,3 +6,4 @@ hiredis
|
|||||||
redis
|
redis
|
||||||
psutil
|
psutil
|
||||||
django-redis
|
django-redis
|
||||||
|
iron-mq
|
||||||
|
|||||||
+4
-1
@@ -10,9 +10,12 @@ django-picklefield==0.3.2
|
|||||||
django-redis==4.2.0
|
django-redis==4.2.0
|
||||||
future==0.15.0
|
future==0.15.0
|
||||||
hiredis==0.2.0
|
hiredis==0.2.0
|
||||||
|
iron-core==1.1.9 # via iron-mq
|
||||||
|
iron-mq==0.7
|
||||||
msgpack-python==0.4.6 # via django-redis
|
msgpack-python==0.4.6 # via django-redis
|
||||||
psutil==3.2.1
|
psutil==3.2.1
|
||||||
python-dateutil==2.4.2 # via arrow
|
python-dateutil==2.4.2 # via arrow, iron-core
|
||||||
redis==2.10.3
|
redis==2.10.3
|
||||||
|
requests==2.7.0 # via iron-core
|
||||||
six==1.9.0 # via python-dateutil
|
six==1.9.0 # via python-dateutil
|
||||||
wcwidth==0.1.4 # via blessed
|
wcwidth==0.1.4 # via blessed
|
||||||
|
|||||||
@@ -26,7 +26,7 @@ class PyTest(Command):
|
|||||||
|
|
||||||
setup(
|
setup(
|
||||||
name='django-q',
|
name='django-q',
|
||||||
version='0.6.2',
|
version='0.6.3',
|
||||||
author='Ilan Steemers',
|
author='Ilan Steemers',
|
||||||
author_email='koed00@gmail.com',
|
author_email='koed00@gmail.com',
|
||||||
keywords='django task queue worker redis disque multiprocessing',
|
keywords='django task queue worker redis disque multiprocessing',
|
||||||
|
|||||||
Reference in New Issue
Block a user