Moves bulk iteration to pusher

This commit is contained in:
Ilan Steemers
2015-09-15 14:45:06 +02:00
parent 9b7c76b9d3
commit a060b3508a
8 changed files with 39 additions and 73 deletions
-1
View File
@@ -7,7 +7,6 @@ class Broker(object):
self.connection = self.get_connection(list_key)
self.list_key = list_key
self.cache = self.get_cache()
self.task_cache = []
def enqueue(self, task):
"""
+3 -11
View File
@@ -17,17 +17,9 @@ class Sqs(Broker):
# sqs supports max 10 messages in bulk
if Conf.BULK > 10:
Conf.BULK = 10
t = None
if len(self.task_cache) > 0:
t = self.task_cache.pop()
else:
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.body
tasks = self.queue.receive_messages(MaxNumberOfMessages=Conf.BULK, VisibilityTimeout=Conf.RETRY)
if tasks:
return [(t.receipt_handle, t.body) for t in tasks]
def acknowledge(self, task_id):
return self.delete(task_id)
+1 -9
View File
@@ -11,18 +11,10 @@ class Disque(Broker):
'ADDJOB {} {} 500 RETRY {}'.format(self.list_key, task, retry)).decode()
def dequeue(self):
t = None
if len(self.task_cache) > 0:
t = self.task_cache.pop()
else:
tasks = self.connection.execute_command(
'GETJOB COUNT {} TIMEOUT 1000 FROM {}'.format(Conf.BULK, self.list_key))
if tasks:
t = tasks.pop()
if tasks:
self.task_cache = tasks
if t:
return t[1].decode(), t[2].decode()
return [(t[1].decode(), t[2].decode()) for t in tasks]
def queue_size(self):
return self.connection.execute_command('QLEN {}'.format(self.list_key))
+1 -9
View File
@@ -9,18 +9,10 @@ class IronMQBroker(Broker):
return self.connection.post(task)['ids'][0]
def dequeue(self):
t = None
if len(self.task_cache) > 0:
t = self.task_cache.pop()
else:
timeout = Conf.RETRY or None
tasks = self.connection.get(timeout=timeout, wait=1, max=Conf.BULK)['messages']
if tasks:
t = tasks.pop()
if tasks:
self.task_cache = tasks
if t:
return t['id'], t['body']
return [(t['id'], t['body']) for t in tasks]
def ping(self):
return self.connection.name == self.list_key
+1 -12
View File
@@ -31,24 +31,13 @@ class ORM(Broker):
return package.pk
def dequeue(self):
if len(self.task_cache) > 0:
t = self.task_cache.pop()
return t.pk, t.payload
else:
# Get new and timed out tasks
tasks = OrmQ.objects.using(Conf.ORM).filter(
Q(key=self.list_key, lock__isnull=True) |
Q(key=self.list_key, lock__lte=timezone.now() - timedelta(seconds=Conf.RETRY)))[:Conf.BULK]
if tasks:
# lock them
OrmQ.objects.using(Conf.ORM).filter(pk__in=tasks).update(lock=timezone.now())
tasks = [t for t in tasks]
# pop one task
t = tasks.pop()
if tasks:
# add remainder to cache
self.task_cache = [t for t in tasks]
return t.pk, t.payload
return [(t.pk, t.payload) for t in tasks]
# empty queue, spare the cpu
sleep(0.2)
+1 -1
View File
@@ -19,7 +19,7 @@ class Redis(Broker):
def dequeue(self):
task = self.connection.blpop(self.list_key, 1)
if task:
return None, task[1]
return [(None, task[1])]
def queue_size(self):
return self.connection.llen(self.list_key)
+13 -12
View File
@@ -298,23 +298,24 @@ def pusher(task_queue, event, broker=None):
logger.info(_('{} pushing tasks at {}').format(current_process().name, current_process().pid))
while True:
try:
task = broker.dequeue()
task_set = broker.dequeue()
except Exception as e:
logger.error(e)
# broker probably crashed. Let the sentinel handle it.
sleep(10)
break
if task:
ack_id = task[0]
# unpack the task
try:
task = signing.SignedPackage.loads(task[1])
except (TypeError, signing.BadSignature) as e:
logger.error(e)
broker.fail(ack_id)
continue
task['ack_id'] = ack_id
task_queue.put(task)
if task_set:
for task in task_set:
ack_id = task[0]
# unpack the task
try:
task = signing.SignedPackage.loads(task[1])
except (TypeError, signing.BadSignature) as e:
logger.error(e)
broker.fail(ack_id)
continue
task['ack_id'] = ack_id
task_queue.put(task)
logger.debug(_('queueing from {}').format(broker.list_key))
if event.is_set():
break
+19 -18
View File
@@ -49,7 +49,7 @@ def test_disque():
broker.enqueue('test')
assert broker.queue_size() == 1
# dequeue
task = broker.dequeue()
task = broker.dequeue()[0]
assert task[1] == 'test'
broker.acknowledge(task[0])
assert broker.queue_size() == 0
@@ -61,7 +61,7 @@ def test_disque():
assert broker.queue_size() == 0
sleep(1.5)
assert broker.queue_size() == 1
task = broker.dequeue()
task = broker.dequeue()[0]
assert broker.queue_size() == 0
broker.acknowledge(task[0])
sleep(1.5)
@@ -78,8 +78,8 @@ def test_disque():
broker.enqueue('test')
Conf.BULK = 5
Conf.DISQUE_FASTACK = True
for i in range(5):
task = broker.dequeue()
tasks = broker.dequeue()
for task in tasks:
assert task is not None
broker.acknowledge(task[0])
# test duplicate acknowledge
@@ -117,7 +117,7 @@ def test_ironmq():
# enqueue
broker.enqueue('test')
# dequeue
task = broker.dequeue()
task = broker.dequeue()[0]
assert task[1] == 'test'
broker.acknowledge(task[0])
assert broker.dequeue() is None
@@ -126,7 +126,7 @@ def test_ironmq():
broker.enqueue('test')
assert broker.dequeue() is not None
sleep(1.5)
task = broker.dequeue()
task = broker.dequeue()[0]
assert len(task) > 0
broker.acknowledge(task[0])
sleep(1.5)
@@ -141,8 +141,8 @@ def test_ironmq():
for i in range(5):
broker.enqueue('test')
Conf.BULK = 5
for i in range(5):
task = broker.dequeue()
tasks = broker.dequeue()
for task in tasks:
assert task is not None
broker.acknowledge(task[0])
# duplicate acknowledge
@@ -174,7 +174,7 @@ def test_sqs():
# enqueue
broker.enqueue('test')
# dequeue
task = broker.dequeue()
task = broker.dequeue()[0]
assert task[1] == 'test'
broker.acknowledge(task[0])
assert broker.dequeue() is None
@@ -183,26 +183,26 @@ def test_sqs():
broker.enqueue('test')
assert broker.dequeue() is not None
sleep(1.5)
task = broker.dequeue()
task = broker.dequeue()[0]
assert len(task) > 0
broker.acknowledge(task[0])
sleep(1.5)
# delete job
broker.enqueue('test')
task_id = broker.dequeue()[0]
task_id = broker.dequeue()[0][0]
broker.delete(task_id)
assert broker.dequeue() is None
# fail
broker.enqueue('test')
while task is None:
task = broker.dequeue()
task = broker.dequeue()[0]
broker.fail(task[0])
# bulk test
for i in range(10):
broker.enqueue('test')
Conf.BULK = 12
for i in range(10):
task = broker.dequeue()
tasks = broker.dequeue()
for task in tasks:
assert task is not None
broker.acknowledge(task[0])
# duplicate acknowledge
@@ -216,6 +216,7 @@ def test_sqs():
Conf.BULK = 1
Conf.DJANGO_REDIS = 'default'
@pytest.mark.django_db
def test_orm():
Conf.ORM = 'default'
@@ -229,7 +230,7 @@ def test_orm():
broker.enqueue('test')
assert broker.queue_size() == 1
# dequeue
task = broker.dequeue()
task = broker.dequeue()[0]
assert task[1] == 'test'
broker.acknowledge(task[0])
assert broker.queue_size() == 0
@@ -241,7 +242,7 @@ def test_orm():
assert broker.queue_size() == 0
sleep(1.5)
assert broker.queue_size() == 1
task = broker.dequeue()
task = broker.dequeue()[0]
assert broker.queue_size() == 0
broker.acknowledge(task[0])
sleep(1.5)
@@ -257,8 +258,8 @@ def test_orm():
for i in range(5):
broker.enqueue('test')
Conf.BULK = 5
for i in range(5):
task = broker.dequeue()
tasks = broker.dequeue()
for task in tasks:
assert task is not None
broker.acknowledge(task[0])
# test duplicate acknowledge