fcbfba63a6
Apply piskvorky/gensim#3649 for the forkserver multiprocessing default.
150 lines
6.3 KiB
Diff
150 lines
6.3 KiB
Diff
diff --git a/gensim/utils.py b/gensim/utils.py
|
|
index 2abd5179e09671de1ea9ef48a9dc3fc4bce9b465..20bad4857bf4e159bc0e244522aa21bd41af8712 100644
|
|
--- a/gensim/utils.py
|
|
+++ b/gensim/utils.py
|
|
@@ -1302,90 +1302,68 @@
|
|
self.q.put(wrapped_chunk.pop(), block=True)
|
|
|
|
|
|
-# Multiprocessing on Windows (and on OSX with python3.8+) uses "spawn" mode, which
|
|
-# causes issues with pickling.
|
|
-# So for these two platforms, use simpler serial processing in `chunkize`.
|
|
-# See https://github.com/RaRe-Technologies/gensim/pull/2800#discussion_r410890171
|
|
-if os.name == 'nt' or (sys.platform == "darwin" and sys.version_info >= (3, 8)):
|
|
- def chunkize(corpus, chunksize, maxsize=0, as_numpy=False):
|
|
- """Split `corpus` into fixed-sized chunks, using :func:`~gensim.utils.chunkize_serial`.
|
|
+def chunkize(corpus, chunksize, maxsize=0, as_numpy=False):
|
|
+ """Split `corpus` into fixed-sized chunks, using :func:`~gensim.utils.chunkize_serial`.
|
|
|
|
- Parameters
|
|
- ----------
|
|
- corpus : iterable of object
|
|
- An iterable.
|
|
- chunksize : int
|
|
- Split `corpus` into chunks of this size.
|
|
- maxsize : int, optional
|
|
- Ignored. For interface compatibility only.
|
|
- as_numpy : bool, optional
|
|
- Yield chunks as `np.ndarray` s instead of lists?
|
|
-
|
|
- Yields
|
|
- ------
|
|
- list OR np.ndarray
|
|
- "chunksize"-ed chunks of elements from `corpus`.
|
|
-
|
|
- """
|
|
- if maxsize > 0:
|
|
- entity = "Windows" if os.name == 'nt' else "OSX with python3.8+"
|
|
- warnings.warn("detected %s; aliasing chunkize to chunkize_serial" % entity)
|
|
- for chunk in chunkize_serial(corpus, chunksize, as_numpy=as_numpy):
|
|
- yield chunk
|
|
-else:
|
|
- def chunkize(corpus, chunksize, maxsize=0, as_numpy=False):
|
|
- """Split `corpus` into fixed-sized chunks, using :func:`~gensim.utils.chunkize_serial`.
|
|
-
|
|
- Parameters
|
|
- ----------
|
|
- corpus : iterable of object
|
|
- An iterable.
|
|
- chunksize : int
|
|
- Split `corpus` into chunks of this size.
|
|
- maxsize : int, optional
|
|
- If > 0, prepare chunks in a background process, filling a chunk queue of size at most `maxsize`.
|
|
- as_numpy : bool, optional
|
|
- Yield chunks as `np.ndarray` instead of lists?
|
|
-
|
|
- Yields
|
|
- ------
|
|
- list OR np.ndarray
|
|
- "chunksize"-ed chunks of elements from `corpus`.
|
|
+ Parameters
|
|
+ ----------
|
|
+ corpus : iterable of object
|
|
+ An iterable.
|
|
+ chunksize : int
|
|
+ Split `corpus` into chunks of this size.
|
|
+ maxsize : int, optional
|
|
+ If > 0, prepare chunks in a background process, filling a chunk queue of size at most `maxsize`.
|
|
+ as_numpy : bool, optional
|
|
+ Yield chunks as `np.ndarray` instead of lists?
|
|
|
|
- Notes
|
|
- -----
|
|
- Each chunk is of length `chunksize`, except the last one which may be smaller.
|
|
- A once-only input stream (`corpus` from a generator) is ok, chunking is done efficiently via itertools.
|
|
+ Yields
|
|
+ ------
|
|
+ list OR np.ndarray
|
|
+ "chunksize"-ed chunks of elements from `corpus`.
|
|
|
|
- If `maxsize > 0`, don't wait idly in between successive chunk `yields`, but rather keep filling a short queue
|
|
- (of size at most `maxsize`) with forthcoming chunks in advance. This is realized by starting a separate process,
|
|
- and is meant to reduce I/O delays, which can be significant when `corpus` comes from a slow medium
|
|
- like HDD, database or network.
|
|
+ Notes
|
|
+ -----
|
|
+ Each chunk is of length `chunksize`, except the last one which may be smaller.
|
|
+ A once-only input stream (`corpus` from a generator) is ok, chunking is done efficiently via itertools.
|
|
|
|
- If `maxsize == 0`, don't fool around with parallelism and simply yield the chunksize
|
|
- via :func:`~gensim.utils.chunkize_serial` (no I/O optimizations).
|
|
+ If `maxsize > 0`, don't wait idly in between successive chunk `yields`, but rather keep filling a short queue
|
|
+ (of size at most `maxsize`) with forthcoming chunks in advance. This is realized by starting a separate process,
|
|
+ and is meant to reduce I/O delays, which can be significant when `corpus` comes from a slow medium
|
|
+ like HDD, database or network.
|
|
|
|
- Yields
|
|
- ------
|
|
- list of object OR np.ndarray
|
|
- Groups based on `iterable`
|
|
+ If `maxsize == 0`, don't fool around with parallelism and simply yield the chunksize
|
|
+ via :func:`~gensim.utils.chunkize_serial` (no I/O optimizations).
|
|
|
|
- """
|
|
- assert chunksize > 0
|
|
+ Yields
|
|
+ ------
|
|
+ list of object OR np.ndarray
|
|
+ Groups based on `iterable`
|
|
|
|
+ """
|
|
+ assert chunksize > 0
|
|
+
|
|
+ # 1. Evaluate the start method dynamically at runtime
|
|
+ # 2. Only attempt multiprocessing if maxsize > 0 AND the method is 'fork'
|
|
+ # See https://github.com/RaRe-Technologies/gensim/pull/2800#discussion_r410890171
|
|
+ if maxsize > 0 and multiprocessing.get_start_method() == "fork":
|
|
+ q = multiprocessing.Queue(maxsize=maxsize)
|
|
+ worker = InputQueue(q, corpus, chunksize, maxsize=maxsize, as_numpy=as_numpy)
|
|
+ worker.daemon = True
|
|
+ worker.start()
|
|
+ while True:
|
|
+ chunk = [q.get(block=True)]
|
|
+ if chunk[0] is None:
|
|
+ break
|
|
+ yield chunk.pop()
|
|
+ else:
|
|
+ # Fallback to serial processing
|
|
if maxsize > 0:
|
|
- q = multiprocessing.Queue(maxsize=maxsize)
|
|
- worker = InputQueue(q, corpus, chunksize, maxsize=maxsize, as_numpy=as_numpy)
|
|
- worker.daemon = True
|
|
- worker.start()
|
|
- while True:
|
|
- chunk = [q.get(block=True)]
|
|
- if chunk[0] is None:
|
|
- break
|
|
- yield chunk.pop()
|
|
- else:
|
|
- for chunk in chunkize_serial(corpus, chunksize, as_numpy=as_numpy):
|
|
- yield chunk
|
|
+ # We only warn if they asked for multiprocessing but we have to deny it
|
|
+ current_method = multiprocessing.get_start_method()
|
|
+ warnings.warn(f"start method is '{current_method}', not 'fork'; aliasing chunkize to chunkize_serial")
|
|
+
|
|
+ for chunk in chunkize_serial(corpus, chunksize, as_numpy=as_numpy):
|
|
+ yield chunk
|
|
|
|
|
|
def smart_extension(fname, ext):
|