From e7e2329660a3f449e5c2f33c853afa9bfe70086e Mon Sep 17 00:00:00 2001 From: Dylan Horkin Date: Fri, 15 Nov 2019 11:26:37 -0800 Subject: [PATCH 1/4] Bugfix: context not provided to init for multiprocessing queue subclass --- util/queue.py | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/util/queue.py b/util/queue.py index 9765658..7161c43 100644 --- a/util/queue.py +++ b/util/queue.py @@ -66,7 +66,7 @@ class Queue(multiprocessing.queues.Queue): """ def __init__(self, *args, **kwargs): - super(Queue, self).__init__(*args, **kwargs) + super(Queue, self).__init__(*args, **kwargs, ctx=multiprocessing.get_context()) self._size = SharedCounter(0) def put(self, *args, **kwargs): From 23906561b361a3798d51d3da92c8ef2c72481431 Mon Sep 17 00:00:00 2001 From: Dylan Horkin Date: Fri, 15 Nov 2019 11:28:37 -0800 Subject: [PATCH 2/4] Bugfix: util.queue._size was not being pickled properly Solved by adding __getstate__ and __setstate__ methods --- util/queue.py | 8 ++++++++ 1 file changed, 8 insertions(+) diff --git a/util/queue.py b/util/queue.py index 7161c43..a975a71 100644 --- a/util/queue.py +++ b/util/queue.py @@ -69,6 +69,14 @@ def __init__(self, *args, **kwargs): super(Queue, self).__init__(*args, **kwargs, ctx=multiprocessing.get_context()) self._size = SharedCounter(0) + # __getstate__ and __setstate__ are needed for pickling, otherwise _size won't be copied. + def __getstate__(self): + return super(Queue, self).__getstate__() + (self._size,) + + def __setstate__(self, state): + self._size = state[-1] + super(Queue, self).__setstate__(state[:-1]) + def put(self, *args, **kwargs): super(Queue, self).put(*args, **kwargs) self._size.increment(1) From ec5faa36cab80b7c419c77f64ec775fba7eda651 Mon Sep 17 00:00:00 2001 From: Dylan Horkin Date: Fri, 27 Dec 2019 13:02:58 -0800 Subject: [PATCH 3/4] Respond to PR feedback on util/queue.py: - fix legacy python compatibility - use namedtuple to pass state information --- util/queue.py | 15 +++++++++++---- 1 file changed, 11 insertions(+), 4 deletions(-) diff --git a/util/queue.py b/util/queue.py index a975a71..62e9922 100644 --- a/util/queue.py +++ b/util/queue.py @@ -20,6 +20,7 @@ import multiprocessing import multiprocessing.queues +from collections import namedtuple class SharedCounter(object): @@ -51,6 +52,9 @@ def value(self): return self.count.value +QueueState = namedtuple('QueueState', ['queue', 'size']) + + class Queue(multiprocessing.queues.Queue): """ A portable implementation of multiprocessing.Queue. @@ -66,16 +70,19 @@ class Queue(multiprocessing.queues.Queue): """ def __init__(self, *args, **kwargs): - super(Queue, self).__init__(*args, **kwargs, ctx=multiprocessing.get_context()) + if sys.version_info >= (3, 4) and 'ctx' not in kwargs: + kwargs['ctx'] = multiprocessing.get_context() + super(Queue, self).__init__(*args, **kwargs) self._size = SharedCounter(0) # __getstate__ and __setstate__ are needed for pickling, otherwise _size won't be copied. def __getstate__(self): - return super(Queue, self).__getstate__() + (self._size,) + return QueueState(queue=super(Queue, self).__getstate__(), + size=self._size) def __setstate__(self, state): - self._size = state[-1] - super(Queue, self).__setstate__(state[:-1]) + self._size = state.size + super(Queue, self).__setstate__(state.queue) def put(self, *args, **kwargs): super(Queue, self).put(*args, **kwargs) From 401b83d7014daed807226347d4fcd21faafe58af Mon Sep 17 00:00:00 2001 From: Dylan Horkin Date: Mon, 6 Jan 2020 16:21:41 -0800 Subject: [PATCH 4/4] Fix missing sys import --- util/queue.py | 1 + 1 file changed, 1 insertion(+) diff --git a/util/queue.py b/util/queue.py index 62e9922..9e4f39d 100644 --- a/util/queue.py +++ b/util/queue.py @@ -18,6 +18,7 @@ # along with this program. If not, see . +import sys import multiprocessing import multiprocessing.queues from collections import namedtuple