Python multiprocessing.Queue无法写入全部值 仅加print可正常运行
问题根因
multiprocessing.Queue 底层依赖「进程内缓冲 + 后台feeder线程」实现:你调用put()写入的数据不会直接发送给对端进程,而是先存到当前进程的内部缓冲,由后台feeder线程异步把缓冲数据写到跨进程Pipe中。
默认情况下,进程退出时会自动join这个feeder线程,等所有缓冲数据都发送完成再退出。你在CODEBLOCK1中调用了metricsQueue.cancel_join_thread(),相当于取消了这个等待逻辑:进程1执行完代码后会直接退出,不会等feeder线程把缓冲里的数据发完,这时候你最后写入的None大概率还没来得及发到Pipe中,直接被丢弃,导致进程2永远收不到终止信号。
现象解释
你观察到的几种正常运行的情况,本质都是给feeder线程留了足够的时间发完所有数据:
- 注释CODEBLOCK1:恢复默认的join逻辑,进程1退出前一定会等feeder把包括
None在内的所有缓冲数据发完,不会丢数据。 - 注释CODEBLOCK2:没有992条数据的写入压力,仅写入的单个
None会被feeder线程立刻发出去,即使取消了join逻辑,进程1退出前数据已经传输完成。 - 添加
print()后正常运行:print是IO操作,会阻塞进程1的执行线程一小段时间,这段时间足够feeder线程把所有缓冲数据写完,None成功发到对端。
你尝试的block=True参数不生效是正常的:该参数仅控制put()时队列满的等待逻辑,你写入时队列一直有空闲,数据已经成功进入了进程内缓冲,block管不到feeder线程的异步传输逻辑。
稳定解决方案
方案1(最推荐)
删除cancel_join_thread()调用,使用默认行为即可。如果你的业务确实有强制杀进程时不能阻塞的需求,可以调整调用时机:在所有业务数据+终止信号None都put完成后,再调用cancel_join_thread(),不要一开始就调用。
方案2
如果必须保留cancel_join_thread()的提前调用,在写入None之后主动等待数据传输完成:
metricsQueue.put(None, block=False) metricsQueue.close() metricsQueue.join_thread() # 主动等feeder线程把所有数据发完
优化建议
你当前进程2的逻辑是忙等:拿不到数据时会无限空转占用CPU,可以给get加一个短超时降低CPU占用:
newMetricsPoint = metricsQueue.get(block=True, timeout=0.1)
内容的提问来源于stack exchange,提问作者agjc

