mirror of
https://github.com/django-q2/django-q2.git
synced 2026-10-08 03:28:12 +08:00
Adds bulk get to ironmq
This commit is contained in:
@@ -7,6 +7,7 @@ class Broker(object):
|
|||||||
self.connection = self.get_connection(list_key)
|
self.connection = self.get_connection(list_key)
|
||||||
self.list_key = list_key
|
self.list_key = list_key
|
||||||
self.cache = self.get_cache()
|
self.cache = self.get_cache()
|
||||||
|
self.task_cache = []
|
||||||
|
|
||||||
def enqueue(self, task):
|
def enqueue(self, task):
|
||||||
"""
|
"""
|
||||||
|
|||||||
@@ -1,18 +1,25 @@
|
|||||||
from django_q.conf import Conf
|
from django_q.conf import Conf, logger
|
||||||
from django_q.brokers import Broker
|
from django_q.brokers import Broker
|
||||||
from iron_mq import IronMQ
|
from iron_mq import IronMQ
|
||||||
|
|
||||||
|
|
||||||
class IronMQBroker(Broker):
|
class IronMQBroker(Broker):
|
||||||
|
|
||||||
def enqueue(self, task):
|
def enqueue(self, task):
|
||||||
return self.connection.post(task)['ids'][0]
|
return self.connection.post(task)['ids'][0]
|
||||||
|
|
||||||
def dequeue(self):
|
def dequeue(self):
|
||||||
timeout = Conf.RETRY or None
|
t = None
|
||||||
task = self.connection.get(timeout=timeout, wait=1)['messages']
|
if len(self.task_cache) > 0:
|
||||||
if task:
|
t = self.task_cache.pop()
|
||||||
return task[0]['id'], task[0]['body']
|
else:
|
||||||
|
timeout = Conf.RETRY or None
|
||||||
|
tasks = self.connection.get(timeout=timeout, wait=1, max=Conf.IRON_MQ_MAX)['messages']
|
||||||
|
if tasks:
|
||||||
|
t = tasks.pop()
|
||||||
|
if tasks:
|
||||||
|
self.task_cache = tasks
|
||||||
|
if t:
|
||||||
|
return t['id'], t['body']
|
||||||
|
|
||||||
def ping(self):
|
def ping(self):
|
||||||
return self.connection.name == self.list_key
|
return self.connection.name == self.list_key
|
||||||
@@ -41,4 +48,4 @@ class IronMQBroker(Broker):
|
|||||||
@staticmethod
|
@staticmethod
|
||||||
def get_connection(list_key=Conf.PREFIX):
|
def get_connection(list_key=Conf.PREFIX):
|
||||||
ironmq = IronMQ(name=None, **Conf.IRON_MQ)
|
ironmq = IronMQ(name=None, **Conf.IRON_MQ)
|
||||||
return ironmq.queue(queue_name=list_key)
|
return ironmq.queue(queue_name=list_key)
|
||||||
|
|||||||
@@ -39,6 +39,7 @@ class Conf(object):
|
|||||||
|
|
||||||
# IronMQ broker
|
# IronMQ broker
|
||||||
IRON_MQ = conf.get('iron_mq', None)
|
IRON_MQ = conf.get('iron_mq', None)
|
||||||
|
IRON_MQ_MAX = conf.get('iron_mq_max', 1)
|
||||||
|
|
||||||
# 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')
|
||||||
|
|||||||
@@ -129,3 +129,4 @@ def test_ironmq():
|
|||||||
assert broker.queue_size() == 0
|
assert broker.queue_size() == 0
|
||||||
# back to django-redis
|
# back to django-redis
|
||||||
Conf.IRON_MQ = None
|
Conf.IRON_MQ = None
|
||||||
|
Conf.DJANGO_REDIS = 'default'
|
||||||
|
|||||||
Reference in New Issue
Block a user