diff --git a/django_q/brokers/aws_sqs.py b/django_q/brokers/aws_sqs.py index 1caa13a..db25215 100644 --- a/django_q/brokers/aws_sqs.py +++ b/django_q/brokers/aws_sqs.py @@ -1,19 +1,17 @@ from django_q.conf import Conf from django_q.brokers import Broker -import boto.sqs -from boto.sqs.message import RawMessage +from boto3 import Session class Sqs(Broker): def __init__(self, list_key=Conf.PREFIX): + self.sqs = None 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 + response = self.queue.send_message(MessageBody=task) + return response.get('MessageId') def dequeue(self): # sqs supports max 10 messages in bulk @@ -23,50 +21,45 @@ class Sqs(Broker): 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) + tasks = self.queue.receive_messages(MaxNumberOfMessages=Conf.BULK, VisibilityTimeout=Conf.RETRY) if tasks: t = tasks.pop() if tasks: self.task_cache = tasks if t: - return t.receipt_handle, t.get_body() + return t.receipt_handle, t.body def acknowledge(self, task_id): return self.delete(task_id) def queue_size(self): - return self.queue.count() + return int(self.queue.attributes['ApproximateNumberOfMessages']) def delete(self, task_id): - m = RawMessage() - m.receipt_handle = task_id - return self.queue.delete_message(m) + message = self.sqs.Message(self.queue.url, task_id) + message.delete() def fail(self, task_id): self.delete(task_id) def delete_queue(self): - self.connection.delete_queue(self.queue) + self.queue.delete() def purge_queue(self): self.queue.purge() def ping(self): - try: - self.connection.get_all_queues() - return True - except Exception as e: - raise e + return 'sqs' in self.connection.get_available_resources() 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 + return Session(aws_access_key_id=Conf.SQS['aws_access_key_id'], + aws_secret_access_key=Conf.SQS['aws_secret_access_key'], + region_name=Conf.SQS['aws_region']) def get_queue(self): - return self.connection.create_queue(self.list_key) + self.sqs = self.connection.resource('sqs') + return self.sqs.create_queue(QueueName=self.list_key) diff --git a/requirements.in b/requirements.in index 956f3c1..476c4a9 100644 --- a/requirements.in +++ b/requirements.in @@ -7,4 +7,4 @@ redis psutil django-redis iron-mq -boto +boto3 diff --git a/requirements.txt b/requirements.txt index ffb54c0..daa05ec 100644 --- a/requirements.txt +++ b/requirements.txt @@ -6,16 +6,20 @@ # arrow==0.6.0 blessed==1.9.5 -boto==2.38.0 +boto3==1.1.3 +botocore==1.2.0 # via boto3 django-picklefield==0.3.2 django-redis==4.2.0 +docutils==0.12 # via botocore future==0.15.1 +futures==2.2.0 # via boto3 hiredis==0.2.0 iron-core==1.1.9 # via iron-mq iron-mq==0.7 +jmespath==0.7.1 # via boto3, botocore msgpack-python==0.4.6 # via django-redis psutil==3.2.1 -python-dateutil==2.4.2 # via arrow, iron-core +python-dateutil==2.4.2 # via arrow, botocore, iron-core redis==2.10.3 requests==2.7.0 # via iron-core six==1.9.0 # via python-dateutil