From 8b3d6b7d2a9a7f37572fb7638a22d8495a7e2212 Mon Sep 17 00:00:00 2001 From: ihuk Date: Wed, 14 Oct 2020 22:21:40 +0200 Subject: [PATCH] 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)