From d7ce945f0374b16137655c4c31125b52aade483a Mon Sep 17 00:00:00 2001 From: Ilan Steemers Date: Wed, 9 Sep 2015 18:10:03 +0200 Subject: [PATCH] Amazon SQS broker --- django_q/brokers/__init__.py | 3 ++ django_q/brokers/aws_sqs.py | 69 ++++++++++++++++++++++++++++++++++ django_q/conf.py | 3 ++ django_q/tests/test_brokers.py | 55 +++++++++++++++++++++++++++ requirements.in | 1 + requirements.txt | 1 + 6 files changed, 132 insertions(+) create mode 100644 django_q/brokers/aws_sqs.py diff --git a/django_q/brokers/__init__.py b/django_q/brokers/__init__.py index 626717c..fed4590 100644 --- a/django_q/brokers/__init__.py +++ b/django_q/brokers/__init__.py @@ -157,6 +157,9 @@ def get_broker(list_key=Conf.PREFIX): elif Conf.IRON_MQ: from brokers import ironmq return ironmq.IronMQBroker(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..93bc16e --- /dev/null +++ b/django_q/brokers/aws_sqs.py @@ -0,0 +1,69 @@ +from django_q.conf import Conf +from django_q.brokers import Broker +import boto.sqs +from boto.sqs.message import RawMessage + + +class Sqs(Broker): + def __init__(self, list_key=Conf.PREFIX): + super(Sqs, self).__init__(list_key) + self.queue = self.get_queue() + + def enqueue(self, task): + m = RawMessage() + m.set_body(task) + self.queue.write(m) + return m.id + + def dequeue(self): + t = None + if len(self.task_cache) > 0: + t = self.task_cache.pop() + else: + tasks = self.queue.get_messages(num_messages=Conf.BULK, visibility_timeout=Conf.RETRY) + if tasks: + t = tasks.pop() + if tasks: + self.task_cache = tasks + if t: + return t.receipt_handle, t.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 = RawMessage() + m.receipt_handle = task_id + return self.queue.delete_message(m) + + def fail(self, task_id): + self.delete(task_id) + + 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 + + def info(self): + return 'AWS SQS' + + @staticmethod + def get_connection(list_key=Conf.PREFIX): + conn = boto.sqs.connect_to_region(Conf.SQS['aws_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/conf.py b/django_q/conf.py index 9463806..c723fac 100644 --- a/django_q/conf.py +++ b/django_q/conf.py @@ -40,6 +40,9 @@ class Conf(object): # IronMQ broker IRON_MQ = conf.get('iron_mq', None) + # SQS broker + SQS = conf.get('sqs', None) + # Name of the cluster or site. For when you run multiple sites on one redis server PREFIX = conf.get('name', 'default') diff --git a/django_q/tests/test_brokers.py b/django_q/tests/test_brokers.py index 3086613..97c6e86 100644 --- a/django_q/tests/test_brokers.py +++ b/django_q/tests/test_brokers.py @@ -148,3 +148,58 @@ def test_ironmq(): # back to django-redis Conf.IRON_MQ = None Conf.DJANGO_REDIS = 'default' + + +@pytest.mark.skipif(not os.getenv('AWS_ACCESS_KEY_ID'), + reason="requires AWS credentials") +def test_sqs(): + Conf.SQS = {'aws_region': os.getenv('AWS_REGION'), + 'aws_access_key_id': os.getenv('AWS_ACCESS_KEY_ID'), + 'aws_secret_access_key': os.getenv('AWS_SECRET_ACCESS_KEY')} + # check broker + broker = get_broker(list_key=uuid()[0]) + assert broker.ping() is True + assert broker.info() is not None + assert broker.queue_size() == 0 + # enqueue + broker.enqueue('test') + # dequeue + task = broker.dequeue() + assert task[1] == 'test' + broker.acknowledge(task[0]) + assert broker.dequeue() is None + # Retry test + Conf.RETRY = 1 + broker.enqueue('test') + 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 + broker.enqueue('test') + task_id = broker.dequeue()[0] + broker.delete(task_id) + assert broker.dequeue() is None + # fail + broker.enqueue('test') + task_id = broker.dequeue()[0] + 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]) + # delete queue + broker.enqueue('test') + broker.enqueue('test') + broker.purge_queue() + assert broker.dequeue() is None + broker.delete_queue() + # back to django-redis + Conf.SQS = None + Conf.DJANGO_REDIS = 'default' diff --git a/requirements.in b/requirements.in index c302143..956f3c1 100644 --- a/requirements.in +++ b/requirements.in @@ -7,3 +7,4 @@ redis psutil django-redis iron-mq +boto diff --git a/requirements.txt b/requirements.txt index 7b03a84..ffb54c0 100644 --- a/requirements.txt +++ b/requirements.txt @@ -6,6 +6,7 @@ # arrow==0.6.0 blessed==1.9.5 +boto==2.38.0 django-picklefield==0.3.2 django-redis==4.2.0 future==0.15.1