Skip to content
This repository was archived by the owner on Oct 26, 2022. It is now read-only.
Open
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
16 changes: 16 additions & 0 deletions util/queue.py
Original file line number Diff line number Diff line change
Expand Up @@ -18,8 +18,10 @@
# along with this program. If not, see <http://www.gnu.org/licenses/>.


import sys
import multiprocessing
import multiprocessing.queues
from collections import namedtuple


class SharedCounter(object):
Expand Down Expand Up @@ -51,6 +53,9 @@ def value(self):
return self.count.value


QueueState = namedtuple('QueueState', ['queue', 'size'])

Copy link
Owner

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Let's make this class non-public: _QueueState.


class Queue(multiprocessing.queues.Queue):
""" A portable implementation of multiprocessing.Queue.

Expand All @@ -66,9 +71,20 @@ class Queue(multiprocessing.queues.Queue):
"""

def __init__(self, *args, **kwargs):
if sys.version_info >= (3, 4) and 'ctx' not in kwargs:
Copy link
Owner

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Please add a comment here explaining why this is necessary (i.e. why we need multiprocessing.get_context() in Python 3.4+).

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.
Copy link
Owner

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Let's drop this comment and instead add docstrings to the two methods, e.g. "Returns the contents to pickle for the instance" and "Sets the state of the instance upon unpickling".

def __getstate__(self):
return QueueState(queue=super(Queue, self).__getstate__(),
size=self._size)

def __setstate__(self, state):
self._size = state.size
super(Queue, self).__setstate__(state.queue)

def put(self, *args, **kwargs):
super(Queue, self).put(*args, **kwargs)
self._size.increment(1)
Expand Down