多进程池提前终止时因返回数据量较大间歇性崩溃的问题求助
多进程池提前终止时因返回数据量较大间歇性崩溃的问题求助
我写了一个并行分析图片的脚本,遇到了个头疼的问题:当我尝试提前退出并调用pool.terminate()时,如果传给进程池的函数返回的数据量较大,就会出现间歇性崩溃。看起来我肯定是漏掉了某些清理步骤,但谷歌了一圈也没找到相关说明,不知道这是不是multiprocessing池的预期行为?
我本来可以过滤输出了事,但那样我就搞不懂背后的原因了,有没有大佬能帮我解释下这是怎么回事?
复现代码
import multiprocessing from signal import signal, SIGINT def pool_initializer(): signal(SIGINT, lambda signum, frame: print(signum, frame)) def work(data): output_dimensions = 256 return [[42 for x in range(output_dimensions)] for y in range(output_dimensions)] if __name__ == "__main__": with multiprocessing.Pool(processes=int(multiprocessing.cpu_count()/2), initializer=pool_initializer) as pool: for result, preview in zip(pool.imap_unordered(work, range(10)), range(10)): ## 不提前退出的话完全不会崩溃 if True: break ## 在pool.terminate()这里会间歇性崩溃 pool.terminate() pool.join() print("Great success!")
报错信息
Traceback (most recent call last): File "H:\wouldntyouliketoknowwaterboy\test.py", line 36, in <module> pool.terminate() ~~~~~~~~~~~~~~^^ File "C:\Users\wouldntyouliketoknowwaterboy\AppData\Local\Programs\Python\Python313\Lib\multiprocessing\pool.py", line 657, in terminate self._terminate() ~~~~~~~~~~~~~~~^^ File "C:\Users\wouldntyouliketoknowwaterboy\AppData\Local\Programs\Python\Python313\Lib\multiprocessing\util.py", line 216, in __call__ res = self._callback(*self._args, **self._kwargs) File "C:\Users\wouldntyouliketoknowwaterboy\AppData\Local\Programs\Python\Python313\Lib\multiprocessing\pool.py", line 703, in _terminate_pool outqueue.put(None) # sentinel ~~~~~~~~~~~~^^^^^^ File "C:\Users\wouldntyouliketoknowwaterboy\AppData\Local\Programs\Python\Python313\Lib\multiprocessing\queues.py", line 394, in put self._writer.send_bytes(obj) ~~~~~~~~~~~~~~~~~~~~~~~^^^^^ File "C:\Users\wouldntyouliketoknowwaterboy\AppData\Local\Programs\Python\Python313\Lib\multiprocessing\connection.py", line 200, in send_bytes self._send_bytes(m[offset:offset + size]) ~~~~~~~~~~~~~~~~^^^^^^^^^^^^^^^^^^^^^^^^^ File "C:\Users\wouldntyouliketoknowwaterboy\AppData\Local\Programs\Python\Python313\Lib\multiprocessing\connection.py", line 287, in _send_bytes raise ValueError("concurrent send_bytes() calls " "are not supported") ValueError: concurrent send_bytes() calls are not supported
补充调试信息
编辑1:后来发现问题和返回数据的可变大小无关,而是和数据量大小有关,已经更新了最小示例,现在返回固定大小但依然会崩溃。
编辑2:启用调试日志multiprocessing.log_to_stderr(multiprocessing.util.DEBUG)后得到了以下输出:
[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 got sentinel [DEBUG/MainProcess] result handler found thread._state=TERMINATE [DEBUG/MainProcess] task handler sending sentinel to result handler Exception in thread Thread-2 (_handle_tasks): [DEBUG/MainProcess] ensuring that outqueue is not full [DEBUG/MainProcess] joining worker handler [DEBUG/MainProcess] terminating workers [DEBUG/MainProcess] joining task handler [DEBUG/MainProcess] result handler exiting: len(cache)=1, thread._state=TERMINATE Traceback (most recent call last): File "C:\Users\babar\AppData\Local\Programs\Python\Python313\Lib\threading.py", line 1041, in _bootstrap_inner self.run() ~~~~~~~~^^ File "C:\Users\babar\AppData\Local\Programs\Python\Python313\Lib\threading.py", line 992, in run self._target(*self._args, **self._kwargs) ~~~~~~~~~~~~^^^^^^^^^^^^^^^^^^^^^^^^^^^^^ File "C:\Users\babar\AppData\Local\Programs\Python\Python313\Lib\multiprocessing\pool.py", line 562, in _handle_tasks outqueue.put(None) ~~~~~~~~~~~~^^^^^^ File "C:\Users\babar\AppData\Local\Programs\Python\Python313\Lib\multiprocessing\queues.py", line 394, in put self._writer.send_bytes(obj) ~~~~~~~~~~~~~~~~~~~~~~~^^^^^ File "C:\Users\babar\AppData\Local\Programs\Python\Python313\Lib\multiprocessing\connection.py", line 200, in send_bytes self._send_bytes(m[offset:offset + size]) ~~~~~~~~~~~~~~~~^^^^^^^^^^^^^^^^^^^^^^^^^ File "C:\Users\babar\AppData\Local\Programs\Python\Python313\Lib\multiprocessing\connection.py", line 287, in _send_bytes raise ValueError("concurrent send_bytes() calls " "are not supported") ValueError: concurrent send_bytes() calls are not supported [DEBUG/MainProcess] joining result handler [DEBUG/MainProcess] joining pool workers
我翻了Windows上Python3.13.2的源码,发现了这些细节:
- 进程池有3个持续循环的线程,直到检测到池处于终止/完成状态:
_worker_handler、_task_handler和_result_handler - 当主线程调用
pool.terminate()时,从terminating pool到removing tasks from inqueue until task handler finished这些日志都是主线程输出的,同时_worker_handler.state和_task_handler.state会被设置为TERMINATE worker handler exiting是_worker_handler线程输出的,是它检测到状态为TERMINATE后的最后一行日志- 关键的冲突点来了:
task handler got sentinel是_task_handler线程在函数执行过程中输出的- 同时,调用
pool.terminate()的主线程会把_result_handler._state设置为TERMINATE _result_handler检测到自己要终止,输出result handler found thread._state=TERMINATE- 之后
_task_handler执行到下一行,输出task handler sending sentinel to result handler——前面这些步骤都发生在两行调试打印之间... - 最终,
_task_handler线程和调用terminate的主线程会同时尝试调用outqueue.put(None):如果_task_handler先执行,就会出现最开始的报错;如果主线程先执行,就会出现编辑2里的报错
所以目前看来,解决方案要么是不提前退出,要么是减小返回数据的大小,另外应该给Python提交一个bug报告。
备注:内容来源于stack exchange,提问作者Babar Shariff
相关产品推荐
相关产品推荐

