Files
django-q2/django_q/brokers/mongo.py
T
Ilan Steemers 826178c463 Adds a MongoDB broker
* adds mongo broker tests
* adds mongo brokers docs

* broker info is now cached when it needs to connect to a server, to reduce traffic.
2015-09-25 20:03:39 +02:00

65 lines
1.7 KiB
Python

from datetime import timedelta
from time import sleep
from bson import ObjectId
from django.utils import timezone
from pymongo import MongoClient
from django_q.brokers import Broker
from django_q.conf import Conf
def _timeout():
return timezone.now() - timedelta(seconds=Conf.RETRY)
class Mongo(Broker):
def __init__(self, list_key=Conf.PREFIX):
super(Mongo, self).__init__(list_key)
self.collection = self.connection[Conf.MONGO_DB][list_key]
@staticmethod
def get_connection(list_key=Conf.PREFIX):
return MongoClient(**Conf.MONGO)
def queue_size(self):
return self.collection.count({'lock': {'$lte': _timeout()}})
def lock_size(self):
return self.collection.count({'lock': {'$gt': _timeout()}})
def purge_queue(self):
return self.delete_queue()
def ping(self):
return self.info is not None
def info(self):
if not self._info:
self._info = 'MongoDB {}'.format(self.connection.server_info()['version'])
return self._info
def fail(self, task_id):
self.delete(task_id)
def enqueue(self, task):
inserted_id = self.collection.insert_one({'payload': task, 'lock': _timeout()}).inserted_id
return str(inserted_id)
def dequeue(self):
task = self.collection.find_one_and_update({'lock': {'$lte': _timeout()}}, {'$set': {'lock': timezone.now()}})
if task:
return [(str(task['_id']), task['payload'])]
# empty queue, spare the cpu
sleep(0.2)
def delete_queue(self):
return self.collection.drop()
def delete(self, task_id):
self.collection.delete_one({'_id': ObjectId(task_id)})
def acknowledge(self, task_id):
return self.delete(task_id)