Python3 multiprocessing nondeterministically hangs on terminate when returning large objects

Viewed 444

I'm using python multiprocessing in a way that requires returning comparably large objects back to the main thread. My example is here https://pastebin.com/xs0X9wWu (on pastebin, as it contains a somewhat large json blob as this error requires a large object) The main code is as follows:

def generate_mutations(original):
    for node in original:
        yield original

def check(task):
    return task

if __name__ == "__main__":
    logger = log_to_stderr()
    logger.setLevel(logging.DEBUG)

    while True:
        with Pool(2) as pool:
            for result in pool.imap(check, generate_mutations(DATA)):
            #for result in map(check, generate_mutations(DATA)):
                print("Stopping")
                break
            print("Loop is done")
        print("Processing is done")

What happens is that after a few iterations of the infinite loop, the pool will fail to terminate and just hang (I see "Loop is done", but not "Processing is done"). Aborting (Ctrl+C) gives the following trace (I also added in the multiprocessing debug log):

[DEBUG/MainProcess] created semlock with handle 139668948529152
[DEBUG/MainProcess] created semlock with handle 139668948525056
[DEBUG/MainProcess] created semlock with handle 139668948520960
[DEBUG/MainProcess] created semlock with handle 139668948516864
[DEBUG/MainProcess] created semlock with handle 139668948512768
[DEBUG/MainProcess] created semlock with handle 139668948508672
[DEBUG/MainProcess] added worker
[INFO/ForkPoolWorker-5] child process calling self.run()
[DEBUG/MainProcess] added worker
[INFO/ForkPoolWorker-6] child process calling self.run()
Stopping
Loop is done
[DEBUG/MainProcess] terminating pool
[DEBUG/MainProcess] finalizing pool
[DEBUG/MainProcess] helping task handler/workers to finish
[DEBUG/MainProcess] removing tasks from inqueue until task handler finished
[DEBUG/MainProcess] worker handler exiting
[DEBUG/MainProcess] task handler found thread._state != RUN
[DEBUG/MainProcess] joining worker handler
[DEBUG/MainProcess] result handler found thread._state=TERMINATE
[DEBUG/MainProcess] task handler sending sentinel to result handler
[DEBUG/MainProcess] terminating workers
[DEBUG/MainProcess] ensuring that outqueue is not full
[DEBUG/MainProcess] joining task handler
Traceback (most recent call last):
  File "mwe.py", line 23, in <module>
    print("Loop is done")
  File "/usr/lib/python3.8/multiprocessing/pool.py", line 736, in __exit__
    self.terminate()
  File "/usr/lib/python3.8/multiprocessing/pool.py", line 654, in terminate
    self._terminate()
  File "/usr/lib/python3.8/multiprocessing/util.py", line 224, in __call__
    res = self._callback(*self._args, **self._kwargs)
  File "/usr/lib/python3.8/multiprocessing/pool.py", line 721, in _terminate_pool
    result_handler.join()
  File "/usr/lib/python3.8/threading.py", line 1011, in join
    self._wait_for_tstate_lock()
  File "/usr/lib/python3.8/threading.py", line 1027, in _wait_for_tstate_lock
    elif lock.acquire(block, timeout):
KeyboardInterrupt
[INFO/MainProcess] process shutting down
[DEBUG/MainProcess] running all "atexit" finalizers with priority >= 0
[DEBUG/MainProcess] running the remaining "atexit" finalizers

There are some similar issues for nontermination of pools (https://bugs.python.org/issue9205, https://bugs.python.org/issue22393, https://bugs.python.org/issue38084), but they all seem to come from worker processes terminating in an unexpected way. This does not seem to be the case here. During debugging, I noticed that the error seems to become less frequent if the data that is returned from check is smaller. I can reproduce the issue reliably on ubuntu with Python 3.8 and 3.9.

If I use regular map() instead of pool.imap(), there seem to be no issues whatsoever.

So what is the issue here? Is this a bug in multiprocessing? Does this code misuse multiprocessing in a serious way? How would I workaround this issue?

0 Answers
Related