From c8f9f0f41aead3cfc8bb7c2f0699ba6096a35b32 Mon Sep 17 00:00:00 2001 From: Odysseas YIAKOUMIS Date: Sat, 3 Oct 2026 23:03:17 +0400 Subject: [PATCH] Fix Pool.imap(buffersize=...) dropping results after close() --- Doc/library/multiprocessing.rst | 3 +++ Lib/multiprocessing/pool.py | 29 +++++++++++++++++------------ Lib/test/_test_multiprocessing.py | 16 +++++++++++++++- 3 files changed, 35 insertions(+), 13 deletions(-) diff --git a/Doc/library/multiprocessing.rst b/Doc/library/multiprocessing.rst index 43fad57139057cd..6919e200299a513 100644 --- a/Doc/library/multiprocessing.rst +++ b/Doc/library/multiprocessing.rst @@ -2558,6 +2558,9 @@ with the :class:`Pool` class. set *buffersize* at least to the number of processes in pool (to consume *iterable* as you go), or even higher (to prefetch the next ``N=buffersize-processes`` arguments). + Once :meth:`join` is called, *buffersize* is no longer honored: + the rest of the *iterable* is submitted without waiting for results + to be yielded. .. versionchanged:: next Added the *buffersize* parameter. diff --git a/Lib/multiprocessing/pool.py b/Lib/multiprocessing/pool.py index ef9460ac0aa6817..c0844cca9bb874e 100644 --- a/Lib/multiprocessing/pool.py +++ b/Lib/multiprocessing/pool.py @@ -191,10 +191,13 @@ def __init__(self, processes=None, initializer=None, initargs=(), self._setup_queues() self._taskqueue = queue.SimpleQueue() # The _taskqueue_buffersize_semaphores exist to allow calling .release() - # on every active semaphore when the pool is terminating to let task_handler - # wake up to stop. It's a set so that each iterator object can efficiently - # deregister its semaphore when iterator finishes. + # on every active semaphore when the pool is terminating or joined to + # let task_handler wake up. It's a set so that each iterator object can + # efficiently deregister its semaphore when iterator finishes. self._taskqueue_buffersize_semaphores = set() + # Set by join(): from then on the buffersize of imap iterators is + # ignored, since there may be nobody left to consume their results. + self._taskqueue_buffersize_ignored = False # The _change_notifier queue exist to wake up self._handle_workers() # when the cache (self._cache) is empty or when there is a change in # the _state variable of the thread that runs _handle_workers. @@ -402,11 +405,11 @@ def _guarded_task_generation(self, result_job, func, iterable, sema=None): else: enumerated_iter = iter(enumerate(iterable)) while True: - sema.acquire() - if self._state != RUN: - # The pool is closing or terminating; stop submitting - # the still-throttled tasks so the task handler can - # finish instead of blocking here forever. + if not self._taskqueue_buffersize_ignored: + sema.acquire() + if self._state == TERMINATE: + # The pool is terminating; stop submitting tasks so + # the task handler can finish. break try: i, x = next(enumerated_iter) @@ -666,10 +669,6 @@ def close(self): self._state = CLOSE self._worker_handler._state = CLOSE self._change_notifier.put(None) - # Wake any task generator throttled on a buffersize semaphore so - # it observes the CLOSE state and stops submitting. - for sema in list(self._taskqueue_buffersize_semaphores): - sema.release() def terminate(self): util.debug('terminating pool') @@ -682,6 +681,12 @@ def join(self): raise ValueError("Pool is still running") elif self._state not in (CLOSE, TERMINATE): raise ValueError("In unknown state") + # Wake any task generator throttled on a buffersize semaphore and let + # it submit the remaining tasks unthrottled: join() has to wait for + # them, and there may be nobody left to consume the results. + self._taskqueue_buffersize_ignored = True + for sema in list(self._taskqueue_buffersize_semaphores): + sema.release() self._worker_handler.join() self._task_handler.join() self._result_handler.join() diff --git a/Lib/test/_test_multiprocessing.py b/Lib/test/_test_multiprocessing.py index 46ed8843fcd0519..8e611e31fc7cf81 100644 --- a/Lib/test/_test_multiprocessing.py +++ b/Lib/test/_test_multiprocessing.py @@ -3239,7 +3239,7 @@ def test_imap_with_buffersize_close_after_partial_consumption( p = self.Pool(2) method = getattr(p, method_name) it = method(sqr, range(1000), buffersize=2) - next(it) + first = next(it) finished = threading.Event() def finalize(): p.close() @@ -3249,6 +3249,20 @@ def finalize(): t.start() t.join(support.SHORT_TIMEOUT) self.assertTrue(finished.is_set(), "close()/join() deadlocked") + # The tasks which were not submitted yet must not be dropped. + self.assertEqual(sorted([first, *it]), list(map(sqr, range(1000)))) + + @warnings_helper.ignore_fork_in_thread_deprecation_warnings() + @support.subTests('method_name', ("imap", "imap_unordered")) + def test_imap_with_buffersize_consumption_after_close(self, method_name): + # close() must not drop the tasks of a buffersize iterator which + # were not submitted yet. + p = self.Pool(2) + method = getattr(p, method_name) + it = method(sqr, range(100), buffersize=2) + p.close() + self.assertEqual(sorted(it), list(map(sqr, range(100)))) + p.join() @support.subTests('method_name', ("imap", "imap_unordered")) def test_imap_and_imap_unordered_with_buffersize_on_empty_iterable(