Repository navigation
ProcessPoolExecutor shutdown hangs after future cancel was requested #94440
Description
Activity
- addedtype-bugAn unexpected behavior, bug, or errorAn unexpected behavior, bug, or error
on Jun 30, 2022 Submitted a PR with test and fix. It can be easily backported, and I’ll be happy to get on it too once we merge on the main line :)
It took some time to pin point the issue, but once done, the fix is pretty surgical. The hang scenario happens when the shutdown sequence begins when the pending work items contains futures, and the all get canceled before the next wait. The addition makes sure there are always some running futures when we wait post shutdown, or none at all, in which case the children stopping kicks off.
@AlexWaygood do you know which expert in multiprocessing I could tag here? I'm eager to see this resolved in the next versions :)
Hi @brianquinlan,
I'd love your feedback on this issue and my suggested fix #94468 :)I experience the same issue (Python 3.8.16). It would be great if a fix (presumably #94468) got merged and backported! @brianquinlan, with the understanding that maintainers are busy and appreciated, would you have time to consider the PR?
I experience the same issue (Python 3.8.16). It would be great if a fix (presumably #94468) got merged and backported! @brianquinlan, with the understanding that maintainers are busy and appreciated, would you have time to consider the PR?
FYI, 3.8 is end of life (3.9 is also) — it now only receives fixes for security-related issues, and this isn't security-related. If you want bugfixes, you'll have to upgrade to at least Python 3.10, which is the oldest Python version that still has bugfixes backported to it.
I'm not active on Python right now so I would not use my commit powers in any case :-(
That’s good to know, thanks anyway!
I was able to reproduce the issue on 3.10 and 3.11 at the time of opening the issue.
It would be really great to merge the PR. The fix is delicate but is pretty straightforward.What if the way forward to get this merged? @AlexWaygood do you know maybe?
I'm sorry you've had to wait such a long time for a review :( unfortunately we don't really have any multiprocessing experts on the core dev team at the moment, and it's a complex module, so few core devs feel confident reviewing PRs relating to multiprocessing.
I'll ask around the core dev team and see if anybody could take a look.
1 remaining item
I experience the same issue (Python 3.8.16). It would be great if a fix (presumably #94468) got merged and backported! @brianquinlan, with the understanding that maintainers are busy and appreciated, would you have time to consider the PR?
FYI, 3.8 is end of life (3.9 is also) — it now only receives fixes for security-related issues, and this isn't security-related. If you want bugfixes, you'll have to upgrade to at least Python 3.10, which is the oldest Python version that still has bugfixes backported to it.
That is quite valid. Our packages are intended to work for python 3.7-3.11, so the PR still solves part of the problem. I did not experience issues with python>=3.9 (where I can call
shutdown(wait=False, cancel_futures=True)), but I don't know if that call is bug-free or whether I was just lucky. At least I now understand whyshutdown(wait=True)hangs -- I've added atime.sleep(0.001)just before shutdown to prevent the hang and this seems to be a functional work-around for python<3.9.I've added a
time.sleep(0.001)just before shutdown to prevent the hang and this seems to be a functional work-around for python<3.9.Yeah, I’ve been doing exactly the same.
- added a commit that references this issue
on Mar 16, 2023 Thanks for the fix! 3.9 and earlier only receive security fixes at this point.
- added a commit that references this issue
on Mar 17, 2023 I think it is still not working on Python3.10.11.
Pressing Ctrl-C after 1-2 seconds, it still runs all submitted jobs:import concurrent.futures import tqdm import time # fun = lambda x: time.sleep(x) def fun(x): print("Sleeping", x) time.sleep(2*x) return x args_list = list(range(10)) max_workers=2 try: with concurrent.futures.ProcessPoolExecutor(max_workers=max_workers) as executor: futures = [] with tqdm.tqdm(total=len(args_list)) as pbar: for args in args_list: futures.append(executor.submit(fun, args)) results = [] for future in concurrent.futures.as_completed(futures): results.append(future.result()) pbar.update(1) # also not working # with concurrent.futures.ProcessPoolExecutor(max_workers=max_workers) as executor: # results = list(tqdm.tqdm(executor.map(fun, args_list), total=len(args_list))) except KeyboardInterrupt: # executor.shutdown(wait=False) print("Cancelling") time.sleep(0.1) executor.shutdown(wait=False, cancel_futures=True) print("Cancelled")
Metadata
Metadata
Assignees
Labels
Projects
- StatusShow more project fieldsDone
Bug report
With a ProcessPoolExecutor, after submitting and quickly canceling a future, a call to
shutdown(wait=True)would hang indefinitely.This happens pretty much on all platforms and all recent Python versions.
Here is a minimal reproduction:
The first submission gets the executor going and creates its internal
queue_management_thread.The second submission appears to get that thread to loop, enter a wait state, and never receive a wakeup event.
Introducing a tiny sleep between the second submit and its cancel request makes the issue disappear. From my initial observation it looks like something in the way the
queue_management_workerinternal loop is structured doesn't handle this edge case well.Shutting down with
wait=Falsewould return immediately as expected, but thequeue_management_threadwould then die with an unhandledOSError: handle is closedexception.Environment
Additional info
When tested with
pytest-timeoutunder Ubuntu and cpython 3.8.13, these are the tracebacks at the moment of timing out:Details
_____________________________________ test _____________________________________ @pytest.mark.timeout(10) def test(): ppe = concurrent.futures.ProcessPoolExecutor(1) ppe.submit(int).result() ppe.submit(int).cancel() > ppe.shutdown(wait=True) test_reproduce_python_bug.py:14: _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ /opt/hostedtoolcache/Python/3.8.13/x64/lib/python3.8/concurrent/futures/process.py:686: in shutdown self._queue_management_thread.join() /opt/hostedtoolcache/Python/3.8.13/x64/lib/python3.8/threading.py:1011: in join self._wait_for_tstate_lock() _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ _ self = <Thread(QueueManagerThread, started daemon 140003176535808)> block = True, timeout = -1 def _wait_for_tstate_lock(self, block=True, timeout=-1): # Issue #18808: wait for the thread state to be gone. # At the end of the thread's life, after all knowledge of the thread # is removed from C data structures, C code releases our _tstate_lock. # This method passes its arguments to _tstate_lock.acquire(). # If the lock is acquired, the C code is done, and self._stop() is # called. That sets ._is_stopped to True, and ._tstate_lock to None. lock = self._tstate_lock if lock is None: # already determined that the C code is done assert self._is_stopped > elif lock.acquire(block, timeout): E Failed: Timeout >10.0s /opt/hostedtoolcache/Python/3.8.13/x64/lib/python3.8/threading.py:1027: Failed ----------------------------- Captured stderr call ----------------------------- +++++++++++++++++++++++++++++++++++ Timeout ++++++++++++++++++++++++++++++++++++ ~~~~~~~~~~~~~~~~~ Stack of QueueFeederThread (140003159754496) ~~~~~~~~~~~~~~~~~ File "/opt/hostedtoolcache/Python/3.8.13/x64/lib/python3.8/threading.py", line 890, in _bootstrap self._bootstrap_inner() File "/opt/hostedtoolcache/Python/3.8.13/x64/lib/python3.8/threading.py", line 932, in _bootstrap_inner self.run() File "/opt/hostedtoolcache/Python/3.8.13/x64/lib/python3.8/threading.py", line 870, in run self._target(*self._args, **self._kwargs) File "/opt/hostedtoolcache/Python/3.8.13/x64/lib/python3.8/multiprocessing/queues.py", line 227, in _feed nwait() File "/opt/hostedtoolcache/Python/3.8.13/x64/lib/python3.8/threading.py", line 302, in wait waiter.acquire() ~~~~~~~~~~~~~~~~ Stack of QueueManagerThread (140003176535808) ~~~~~~~~~~~~~~~~~ File "/opt/hostedtoolcache/Python/3.8.13/x64/lib/python3.8/threading.py", line 890, in _bootstrap self._bootstrap_inner() File "/opt/hostedtoolcache/Python/3.8.13/x64/lib/python3.8/threading.py", line 932, in _bootstrap_inner self.run() File "/opt/hostedtoolcache/Python/3.8.13/x64/lib/python3.8/threading.py", line 870, in run self._target(*self._args, **self._kwargs) File "/opt/hostedtoolcache/Python/3.8.13/x64/lib/python3.8/concurrent/futures/process.py", line 362, in _queue_management_worker ready = mp.connection.wait(readers + worker_sentinels) File "/opt/hostedtoolcache/Python/3.8.13/x64/lib/python3.8/multiprocessing/connection.py", line 931, in wait ready = selector.select(timeout) File "/opt/hostedtoolcache/Python/3.8.13/x64/lib/python3.8/selectors.py", line 415, in select fd_event_list = self._selector.poll(timeout) +++++++++++++++++++++++++++++++++++ Timeout ++++++++++++++++++++++++++++++++++++Tracebacks in PyPy are similar on the
concurrent.futures.processlevel. Tracebacks in Windows are different in the lower-level areas, but again similar on theconcurrent.futures.processlevel.Linked PRs:
Linked PRs