mirror of
https://github.com/django-q2/django-q2.git
synced 2026-10-06 13:08:11 +08:00
@@ -13,6 +13,14 @@ class Broker:
|
|||||||
self.cache = self.get_cache()
|
self.cache = self.get_cache()
|
||||||
self._info = None
|
self._info = None
|
||||||
|
|
||||||
|
def __getstate__(self):
|
||||||
|
return self.list_key, self._info
|
||||||
|
|
||||||
|
def __setstate__(self, state):
|
||||||
|
self.list_key, self._info = state
|
||||||
|
self.connection = self.get_connection(self.list_key)
|
||||||
|
self.cache = self.get_cache()
|
||||||
|
|
||||||
def enqueue(self, task):
|
def enqueue(self, task):
|
||||||
"""
|
"""
|
||||||
Puts a task onto the queue
|
Puts a task onto the queue
|
||||||
|
|||||||
@@ -10,6 +10,11 @@ class Sqs(Broker):
|
|||||||
super(Sqs, self).__init__(list_key)
|
super(Sqs, self).__init__(list_key)
|
||||||
self.queue = self.get_queue()
|
self.queue = self.get_queue()
|
||||||
|
|
||||||
|
def __setstate__(self, state):
|
||||||
|
super(Sqs, self).__setstate__(state)
|
||||||
|
self.sqs = None
|
||||||
|
self.queue = self.get_queue()
|
||||||
|
|
||||||
def enqueue(self, task):
|
def enqueue(self, task):
|
||||||
response = self.queue.send_message(MessageBody=task)
|
response = self.queue.send_message(MessageBody=task)
|
||||||
return response.get("MessageId")
|
return response.get("MessageId")
|
||||||
|
|||||||
@@ -19,6 +19,10 @@ class Mongo(Broker):
|
|||||||
super(Mongo, self).__init__(list_key)
|
super(Mongo, self).__init__(list_key)
|
||||||
self.collection = self.get_collection()
|
self.collection = self.get_collection()
|
||||||
|
|
||||||
|
def __setstate__(self, state):
|
||||||
|
super(Mongo, self).__setstate__(state)
|
||||||
|
self.collection = self.get_collection()
|
||||||
|
|
||||||
@staticmethod
|
@staticmethod
|
||||||
def get_connection(list_key: str = Conf.PREFIX) -> MongoClient:
|
def get_connection(list_key: str = Conf.PREFIX) -> MongoClient:
|
||||||
return MongoClient(**Conf.MONGO)
|
return MongoClient(**Conf.MONGO)
|
||||||
|
|||||||
+9
-1
@@ -12,7 +12,15 @@ from time import sleep
|
|||||||
import arrow
|
import arrow
|
||||||
|
|
||||||
# Django
|
# Django
|
||||||
from django import db
|
from django import db, core
|
||||||
|
from django.apps.registry import apps
|
||||||
|
|
||||||
|
try:
|
||||||
|
apps.check_apps_ready()
|
||||||
|
except core.exceptions.AppRegistryNotReady:
|
||||||
|
import django
|
||||||
|
django.setup()
|
||||||
|
|
||||||
from django.conf import settings
|
from django.conf import settings
|
||||||
from django.utils import timezone
|
from django.utils import timezone
|
||||||
from django.utils.translation import gettext_lazy as _
|
from django.utils.translation import gettext_lazy as _
|
||||||
|
|||||||
+7
-1
@@ -54,9 +54,15 @@ class Queue(multiprocessing.queues.Queue):
|
|||||||
super(Queue, self).__init__(
|
super(Queue, self).__init__(
|
||||||
*args, ctx=multiprocessing.get_context(), **kwargs
|
*args, ctx=multiprocessing.get_context(), **kwargs
|
||||||
)
|
)
|
||||||
|
|
||||||
self.size = SharedCounter(0)
|
self.size = SharedCounter(0)
|
||||||
|
|
||||||
|
def __getstate__(self):
|
||||||
|
return super(Queue, self).__getstate__() + (self.size, )
|
||||||
|
|
||||||
|
def __setstate__(self, state):
|
||||||
|
super(Queue, self).__setstate__(state[:-1])
|
||||||
|
self.size = state[-1]
|
||||||
|
|
||||||
def put(self, *args, **kwargs):
|
def put(self, *args, **kwargs):
|
||||||
super(Queue, self).put(*args, **kwargs)
|
super(Queue, self).put(*args, **kwargs)
|
||||||
self.size.increment(1)
|
self.size.increment(1)
|
||||||
|
|||||||
Reference in New Issue
Block a user