Skip to content
Open
Show file tree
Hide file tree
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
3 changes: 3 additions & 0 deletions Doc/library/multiprocessing.rst
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down
29 changes: 17 additions & 12 deletions Lib/multiprocessing/pool.py
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down Expand Up @@ -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:

@picnixz picnixz Oct 3, 2026 •

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

Why this change on the state?

Copy link
Copy Markdown
Author

Choose a reason for hiding this comment

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

The old check, self._state != RUN, is what caused the bug. It is true in both the CLOSE and TERMINATE states, so the generator stopped as soon as close() was called and the items not yet submitted were dropped.

close() should only prevent new submissions. The items of an imap() call that was already made should still be processed, as they are without buffersize. So the generator must keep going in the CLOSE state.

TERMINATE is the only state where dropping the remaining items is correct, so that is now the only state that breaks out of the loop.

The hang from gh-155477 is handled in join() instead, which stops the throttling and wakes the generator.

# The pool is terminating; stop submitting tasks so
# the task handler can finish.
break
try:
i, x = next(enumerated_iter)
Expand Down Expand Up @@ -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')
Expand All @@ -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()
Expand Down
16 changes: 15 additions & 1 deletion Lib/test/_test_multiprocessing.py
Original file line number Diff line number Diff line change
Expand Up @@ -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()
Expand All @@ -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(
Expand Down
Loading