From 8b3d6b7d2a9a7f37572fb7638a22d8495a7e2212 Mon Sep 17 00:00:00 2001 From: ihuk Date: Wed, 14 Oct 2020 22:21:40 +0200 Subject: [PATCH 1/2] Initial work on fix for #424 - TypeError: can't pickle _thread.lock objects. Made `Broker` and `Queue` classes pickable which fixes the the problem with new default `spawn` proces start method. --- django_q/brokers/__init__.py | 8 ++++++++ django_q/cluster.py | 10 +++++++++- django_q/queues.py | 8 +++++++- 3 files changed, 24 insertions(+), 2 deletions(-) diff --git a/django_q/brokers/__init__.py b/django_q/brokers/__init__.py index 5bef988..728290f 100644 --- a/django_q/brokers/__init__.py +++ b/django_q/brokers/__init__.py @@ -13,6 +13,14 @@ class Broker: self.cache = self.get_cache() 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): """ Puts a task onto the queue diff --git a/django_q/cluster.py b/django_q/cluster.py index 1d55663..5fc437e 100644 --- a/django_q/cluster.py +++ b/django_q/cluster.py @@ -12,7 +12,15 @@ from time import sleep import arrow # 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.utils import timezone from django.utils.translation import gettext_lazy as _ diff --git a/django_q/queues.py b/django_q/queues.py index 9844fe2..346ccec 100644 --- a/django_q/queues.py +++ b/django_q/queues.py @@ -54,9 +54,15 @@ class Queue(multiprocessing.queues.Queue): super(Queue, self).__init__( *args, ctx=multiprocessing.get_context(), **kwargs ) - 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): super(Queue, self).put(*args, **kwargs) self.size.increment(1) From 507636716856203d096438653419b1fe0054dab5 Mon Sep 17 00:00:00 2001 From: ihuk Date: Fri, 16 Oct 2020 18:40:26 +0200 Subject: [PATCH 2/2] More work on fix for #424 - TypeError: can't pickle _thread.lock objects. AWS and Mongo brockers should be pickeled correctly now. --- django_q/brokers/aws_sqs.py | 5 +++++ django_q/brokers/mongo.py | 4 ++++ 2 files changed, 9 insertions(+) diff --git a/django_q/brokers/aws_sqs.py b/django_q/brokers/aws_sqs.py index 1a6deda..a666934 100644 --- a/django_q/brokers/aws_sqs.py +++ b/django_q/brokers/aws_sqs.py @@ -10,6 +10,11 @@ class Sqs(Broker): super(Sqs, self).__init__(list_key) 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): response = self.queue.send_message(MessageBody=task) return response.get("MessageId") diff --git a/django_q/brokers/mongo.py b/django_q/brokers/mongo.py index a130113..a4857bd 100644 --- a/django_q/brokers/mongo.py +++ b/django_q/brokers/mongo.py @@ -19,6 +19,10 @@ class Mongo(Broker): super(Mongo, self).__init__(list_key) self.collection = self.get_collection() + def __setstate__(self, state): + super(Mongo, self).__setstate__(state) + self.collection = self.get_collection() + @staticmethod def get_connection(list_key: str = Conf.PREFIX) -> MongoClient: return MongoClient(**Conf.MONGO)