You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

多进程池提前终止时因返回数据量较大间歇性崩溃的问题求助

多进程池提前终止时因返回数据量较大间歇性崩溃的问题求助

我写了一个并行分析图片的脚本,遇到了个头疼的问题:当我尝试提前退出并调用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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.04.14 10:29:37