From 8ce441bed5850faddc59a4b333135fb2d83f8551 Mon Sep 17 00:00:00 2001 From: Ilan Steemers Date: Wed, 2 Sep 2015 19:04:03 +0200 Subject: [PATCH] Adds more brokers * added IronMQ * added Aws SQS * removed SafeRedis. During tests the 'safe' part wasn't consistent enough --- django_q/__init__.py | 2 +- django_q/brokers/__init__.py | 94 +++++++++++++++++++++++++++++--- django_q/brokers/aws_sqs.py | 57 +++++++++++++++++++ django_q/brokers/disque.py | 26 ++++++--- django_q/brokers/djangoredis.py | 10 ---- django_q/brokers/iron_mq.py | 38 +++++++++++++ django_q/brokers/redis_broker.py | 20 +++++-- django_q/conf.py | 18 ++++-- django_q/tests/test_brokers.py | 80 +++++++++++++++++++++++++-- docs/conf.py | 4 +- requirements.in | 2 + requirements.txt | 8 ++- setup.py | 4 +- 13 files changed, 314 insertions(+), 49 deletions(-) create mode 100644 django_q/brokers/aws_sqs.py delete mode 100644 django_q/brokers/djangoredis.py create mode 100644 django_q/brokers/iron_mq.py diff --git a/django_q/__init__.py b/django_q/__init__.py index c0913fb..9fc91f1 100644 --- a/django_q/__init__.py +++ b/django_q/__init__.py @@ -9,6 +9,6 @@ from .models import Task, Schedule, Success, Failure from .cluster import Cluster from .status import Stat -VERSION = (0, 5, 3) +VERSION = (0, 6, 0) default_app_config = 'django_q.apps.DjangoQConfig' diff --git a/django_q/brokers/__init__.py b/django_q/brokers/__init__.py index a57ea78..ce7f9a8 100644 --- a/django_q/brokers/__init__.py +++ b/django_q/brokers/__init__.py @@ -3,30 +3,74 @@ from django.core.cache import caches, InvalidCacheBackendError class Broker(object): - def __init__(self, list_key=Conf.Q_LIST): - self.connection = self.get_connection() + def __init__(self, list_key=Conf.PREFIX): + self.connection = self.get_connection(list_key) self.list_key = list_key self.cache = self.get_cache() def enqueue(self, task): + """ + Puts a task onto the queue + :type task: str + :return: task id + """ pass def dequeue(self): + """ + Gets a task from the queue + :return: tuple with task id and task message + """ pass def queue_size(self): + """ + :return: the amount of tasks in the queue + """ pass - def delete_queue(self, list_key=None): + def delete_queue(self): + """ + Deletes the queue from the broker + """ pass - def acknowledge(self, ack_id): + def purge_queue(self): + """ + Purges the queue of any tasks + """ + pass + + def delete(self, task_id): + """ + Deletes a task from the queue + :param task_id: the id of the task + """ + pass + + def acknowledge(self, task_id): + """ + Acknowledges completion of the task and removes it from the queue. + :param task_id: the id of the task + """ pass def ping(self): + """ + Checks whether the broker connection is available + :rtype: bool + """ pass def set_stat(self, key, value, timeout): + """ + Saves a cluster statistic to the cache provider + :type key: str + :type value: str + :type timeout: int + """ + if not self.cache: + return key_list = self.cache.get(Conf.Q_STAT, []) if key not in key_list: key_list.append(key) @@ -34,9 +78,23 @@ class Broker(object): return self.cache.set(key, value, timeout) def get_stat(self, key): + """ + Gets a cluster statistic from the cache provider + :type key: str + :return: a cluster Stat + """ + if not self.cache: + return return self.cache.get(key) def get_stats(self, pattern): + """ + Returns a list of all cluster stats from the cache provider + :type pattern: str + :return: a list of Stats + """ + if not self.cache: + return key_list = self.cache.get(Conf.Q_STAT) if not key_list or len(key_list) == 0: return [] @@ -52,23 +110,41 @@ class Broker(object): @staticmethod def get_cache(): + """ + Gets the current cache provider + :return: a cache provider + """ try: return caches[Conf.CACHE] except InvalidCacheBackendError: return None @staticmethod - def get_connection(): + def get_connection(list_key=Conf.PREFIX): + """ + Gets a connection to the broker + :param list_key: Optional queue name + :return: a broker connection + """ return 0 -def get_broker(list_key=Conf.Q_LIST): - if Conf.DJANGO_REDIS: - from brokers import djangoredis - return djangoredis.DjangoRedis(list_key=list_key) +def get_broker(list_key=Conf.PREFIX): + """ + Gets the configured broker type + :param list_key: optional queue name + :type list_key: str + :return: + """ + if Conf.IRONMQ: + from brokers import iron_mq + return iron_mq.IronMQBroker(list_key=list_key) elif Conf.DISQUE: from brokers import disque return disque.Disque(list_key=list_key) + elif Conf.SQS: + from brokers import aws_sqs + return aws_sqs.Sqs(list_key=list_key) # default to redis else: from brokers import redis_broker diff --git a/django_q/brokers/aws_sqs.py b/django_q/brokers/aws_sqs.py new file mode 100644 index 0000000..4e58bae --- /dev/null +++ b/django_q/brokers/aws_sqs.py @@ -0,0 +1,57 @@ +import os +from django_q.conf import Conf +from django_q.brokers import Broker +import boto.sqs +from boto.sqs.message import Message + + +class Sqs(Broker): + def __init__(self, list_key=Conf.PREFIX): + super().__init__(list_key) + self.queue = self.get_queue() + + def enqueue(self, task): + m = Message() + m.set_body(task) + self.queue.write(m) + return m.id + + def dequeue(self): + rs = self.queue.get_messages(visibility_timeout=Conf.RETRY or 30) + if rs: + m = rs[0] + return m.receipt_handle, m.get_body() + + def acknowledge(self, task_id): + return self.delete(task_id) + + def queue_size(self): + return self.queue.count() + + def delete(self, task_id): + m = Message() + m.receipt_handle = task_id + return self.queue.delete_message(m) + + def delete_queue(self): + self.connection.delete_queue(self.queue) + + def purge_queue(self): + self.queue.purge() + + def ping(self): + try: + self.connection.get_all_queues() + return True + except Exception as e: + raise e + + @staticmethod + def get_connection(list_key=Conf.PREFIX): + conn = boto.sqs.connect_to_region(Conf.SQS['region'], + aws_access_key_id=Conf.SQS['aws_access_key_id'], + aws_secret_access_key=Conf.SQS['aws_secret_access_key']) + return conn + + def get_queue(self): + return self.connection.create_queue(self.list_key) diff --git a/django_q/brokers/disque.py b/django_q/brokers/disque.py index 45f20d5..d26d29b 100644 --- a/django_q/brokers/disque.py +++ b/django_q/brokers/disque.py @@ -1,9 +1,11 @@ +import random import redis from django_q.brokers import Broker from django_q.conf import Conf class Disque(Broker): + def enqueue(self, task): return self.connection.execute_command( 'ADDJOB {} {} 500 RETRY {}'.format(self.list_key, task, Conf.RETRY)).decode() @@ -16,24 +18,34 @@ class Disque(Broker): def queue_size(self): return self.connection.execute_command('QLEN {}'.format(self.list_key)) - def acknowledge(self, ack_id): - return self.connection.execute_command('ACKJOB {}'.format(ack_id)) + def acknowledge(self, task_id): + return self.connection.execute_command('ACKJOB {}'.format(task_id)) def ping(self): return self.connection.ping() - def delete_queue(self, list_key=None): - raise NotImplementedError + def delete(self, task_id): + return self.connection.execute_command('DELJOB {}'.format(task_id)) + + def delete_queue(self): + jobs = self.connection.execute_command('JSCAN QUEUE {}'.format(self.list_key))[1] + if jobs: + self.connection.execute_command('DELJOB {}'.format(' '.join(map(str, jobs)))) @staticmethod - def get_connection(): + def get_connection(list_key=Conf.PREFIX): + # randomize nodes + random.shuffle(Conf.DISQUE) + # find one that works for node in Conf.DISQUE: host, port = node.split(':') redis_client = redis.Redis(host, int(port)) try: - redis_client.ping() + if Conf.DISQUE_AUTH: + redis_client.execute_command('AUTH {}'.format(Conf.DISQUE_AUTH)) redis_client.decode_responses = True + redis_client.execute_command('HELLO') return redis_client except redis.exceptions.ConnectionError: - pass + continue raise ConnectionError('Could not connect to any Disque nodes') diff --git a/django_q/brokers/djangoredis.py b/django_q/brokers/djangoredis.py deleted file mode 100644 index 1b519d8..0000000 --- a/django_q/brokers/djangoredis.py +++ /dev/null @@ -1,10 +0,0 @@ -import django_redis -from django_q.brokers import redis_broker -from django_q.conf import Conf - - -class DjangoRedis(redis_broker.Redis): - - @staticmethod - def get_connection(): - return django_redis.get_redis_connection(Conf.DJANGO_REDIS) \ No newline at end of file diff --git a/django_q/brokers/iron_mq.py b/django_q/brokers/iron_mq.py new file mode 100644 index 0000000..9fbc79d --- /dev/null +++ b/django_q/brokers/iron_mq.py @@ -0,0 +1,38 @@ +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 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()['msg'] + + def delete(self, task_id): + return self.connection.delete(task_id)['msg'] + + def acknowledge(self, task_id): + return self.delete(task_id) + + @staticmethod + def get_connection(list_key=Conf.PREFIX): + ironmq = IronMQ(name=None, **Conf.IRONMQ) + return ironmq.queue(queue_name=list_key) diff --git a/django_q/brokers/redis_broker.py b/django_q/brokers/redis_broker.py index 042a0cd..de85882 100644 --- a/django_q/brokers/redis_broker.py +++ b/django_q/brokers/redis_broker.py @@ -2,8 +2,17 @@ import redis from django_q.brokers import Broker from django_q.conf import Conf, logger +try: + import django_redis +except ImportError: + django_redis = None + class Redis(Broker): + + def __init__(self, list_key=Conf.PREFIX): + super().__init__(list_key='django_q:{}:q'.format(list_key)) + def enqueue(self, task): return self.connection.rpush(self.list_key, task) @@ -15,14 +24,13 @@ class Redis(Broker): def queue_size(self): return self.connection.llen(self.list_key) - def delete_queue(self, list_key=None): - list_key = list_key if list_key else self.list_key - return self.connection.delete(list_key) + def delete_queue(self): + return self.connection.delete(self.list_key) def ping(self): try: return self.connection.ping() - except Exception as e: + except redis.ConnectionError as e: logger.error('Can not connect to Redis server.') raise e @@ -39,5 +47,7 @@ class Redis(Broker): return self.connection.mget(keys) @staticmethod - def get_connection(): + def get_connection(list_key=Conf.PREFIX): + if django_redis and Conf.DJANGO_REDIS: + return django_redis.get_redis_connection(Conf.DJANGO_REDIS) return redis.StrictRedis(**Conf.REDIS) diff --git a/django_q/conf.py b/django_q/conf.py index 394861c..0e9c7b7 100644 --- a/django_q/conf.py +++ b/django_q/conf.py @@ -35,6 +35,16 @@ class Conf(object): # Disque broker DISQUE = conf.get('disque', None) + # Optional Authentication + DISQUE_AUTH = conf.get('disque_auth', None) + + # Amazon SQS broker + SQS = conf.get('sqs', None) + + # IronMQ broker + IRONMQ = conf.get('ironmq', None) + if IRONMQ and os.environ.get('IRONMQ_TOKEN'): + IRONMQ['token'] = os.environ['IRONMQ_TOKEN'] # Name of the cluster or site. For when you run multiple sites on one redis server PREFIX = conf.get('name', 'default') @@ -46,9 +56,6 @@ class Conf(object): # Failures are always saved SAVE_LIMIT = conf.get('save_limit', 250) - # Maximum number of tasks that each cluster can work on - QUEUE_LIMIT = conf.get('queue_limit', None) - # Number of workers in the pool. Default is cpu count if implemented, otherwise 4. WORKERS = conf.get('workers', False) if not WORKERS: @@ -63,6 +70,9 @@ class Conf(object): # sensible default WORKERS = 4 + # Maximum number of tasks that each cluster can work on + QUEUE_LIMIT = conf.get('queue_limit', None) + # Sets compression of redis packages COMPRESSED = conf.get('compress', False) @@ -96,8 +106,6 @@ class Conf(object): # Django itself should raise an error if it's not configured SECRET_KEY = settings.SECRET_KEY - # The redis list key - Q_LIST = 'django_q:{}:q'.format(PREFIX) # The redis stats key Q_STAT = 'django_q:{}:cluster'.format(PREFIX) diff --git a/django_q/tests/test_brokers.py b/django_q/tests/test_brokers.py index 4e9ff11..f666cff 100644 --- a/django_q/tests/test_brokers.py +++ b/django_q/tests/test_brokers.py @@ -1,5 +1,6 @@ from time import sleep import pytest +import os from django_q.conf import Conf from django_q.brokers import get_broker, Broker @@ -30,10 +31,13 @@ def test_redis(): Conf.DJANGO_REDIS = 'default' -def disabled_test_disque(): +@pytest.mark.skipif(Conf.DISQUE_AUTH is None, + reason="No disque server configured") +def test_disque(): Conf.DISQUE = ['127.0.0.1:7711'] - broker = get_broker() + broker = get_broker(list_key='disque_test') assert broker.ping() is True + broker.delete_queue() broker.enqueue('test') assert broker.queue_size() == 1 task = broker.dequeue() @@ -53,12 +57,76 @@ def disabled_test_disque(): broker.acknowledge(task[0]) sleep(1.5) assert broker.queue_size() == 0 - # errors - with pytest.raises(NotImplementedError): - broker.delete_queue() Conf.DISQUE = ['127.0.0.1:7712', '127.0.0.1:7713'] with pytest.raises(ConnectionError): broker.get_connection() + broker.delete_queue() + assert broker.queue_size() == 0 + # back to django-redis + Conf.DISQUE = None + + +@pytest.mark.skipif(not os.getenv('AWS_ACCESS_KEY_ID'), + reason="requires AWS SQS credentials") +def test_sqs(): + Conf.SQS = {'aws_access_key_id': os.getenv('AWS_ACCESS_KEY_ID'), + 'aws_secret_access_key': os.getenv('AWS_SECRET_ACCESS_KEY'), + 'region': os.getenv('SQS_REGION', 'eu-west-1')} + broker = get_broker(list_key='sqs_test') + assert broker.ping() is True + if broker.queue_size() > 0: + broker.purge_queue() + broker.enqueue('test') + task = broker.dequeue() + assert task[1] == 'test' + broker.acknowledge(task[0]) + assert broker.queue_size() == 0 + # Retry test + Conf.RETRY = 1 + broker.enqueue('test') + broker.dequeue() + assert broker.queue_size() == 0 + sleep(2) + # task should re-queue + assert broker.queue_size() == 1 + task = broker.dequeue() + assert task[1] == 'test' + assert broker.acknowledge(task[0]) is True + assert broker.queue_size() == 0 + broker.delete_queue() + # back to defaults + Conf.SQS = None + + +@pytest.mark.skipif(not os.getenv('IRONMQ_TOKEN'), + reason="requires IronMQ credentials") +def test_ironmq(): + Conf.IRONMQ = {'host': 'mq-aws-eu-west-1.iron.io', + 'token': os.getenv('IRONMQ_TOKEN'), + 'project_id': os.getenv('IRONMQ_PROJECT')} + broker = get_broker(list_key='djangoQ_test') + assert broker.ping() is True + broker.delete_queue() + broker.enqueue('test') + assert broker.queue_size() == 1 + 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 + broker.delete_queue() + assert broker.queue_size() == 0 # back to django-redis Conf.DISQUE = None - Conf.DJANGO_REDIS = 'default' diff --git a/docs/conf.py b/docs/conf.py index 4775ee9..4140e5c 100644 --- a/docs/conf.py +++ b/docs/conf.py @@ -70,9 +70,9 @@ author = 'Ilan Steemers' # built documents. # # The short X.Y version. -version = '0.5' +version = '0.6' # The full version, including alpha/beta/rc tags. -release = '0.5.3' +release = '0.6.0' # The language for content autogenerated by Sphinx. Refer to documentation # for a list of supported languages. diff --git a/requirements.in b/requirements.in index 1ab7ca1..0198050 100644 --- a/requirements.in +++ b/requirements.in @@ -6,3 +6,5 @@ hiredis redis psutil django-redis +boto +iron-mq diff --git a/requirements.txt b/requirements.txt index 83e7e99..2ac9cf2 100644 --- a/requirements.txt +++ b/requirements.txt @@ -6,13 +6,17 @@ # arrow==0.6.0 blessed==1.9.5 +boto==2.38.0 django-picklefield==0.3.1 django-redis==4.2.0 future==0.15.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 -psutil==3.1.1 -python-dateutil==2.4.2 # via arrow +psutil==3.2.0 +python-dateutil==2.4.2 # via arrow, iron-core redis==2.10.3 +requests==2.7.0 # via iron-core six==1.9.0 # via django-picklefield, python-dateutil wcwidth==0.1.4 # via blessed diff --git a/setup.py b/setup.py index aad8571..35b0afc 100644 --- a/setup.py +++ b/setup.py @@ -26,7 +26,7 @@ class PyTest(Command): setup( name='django-q', - version='0.5.3', + version='0.6.0', author='Ilan Steemers', author_email='koed00@gmail.com', keywords='django task queue worker redis multiprocessing', @@ -36,7 +36,7 @@ setup( license='MIT', description='A multiprocessing task queue for Django', long_description=README, - install_requires=['django>=1.7', 'redis', 'django-picklefield', 'blessed', 'arrow', 'future'], + install_requires=['django>=1.7', 'django-picklefield', 'blessed', 'arrow', 'future'], test_requires=['pytest', 'pytest-django', ], cmdclass={'test': PyTest}, classifiers=[