Manager Queue跨进程传递前,进程内线程调用get()触发Broken Pipe错误
问题分析与解决方案
问题重现
在Python 3.7 Linux环境下,以下测试用例会稳定触发BrokenPipeError日志(需通过multiprocessing.get_logger()查看,该错误不会主动抛出):
from multiprocessing import Process, Manager from threading import Thread from time import sleep def test_process(): def listener(q): t = Thread(target=lambda x: print(x.get(True)), args=(q, ), daemon=True) t.start() def publisher(q): q.put('test') with Manager() as mp: q = mp.Queue() a = Process(target=listener, args=(q, ), daemon=True) a.start() sleep(2) b = Process(target=publisher, args=(q, ), daemon=True) b.start() a.kill() b.kill()
运行命令:
pytest -k test_process -vv
报错日志:
exception in thread serving 'Process-2|Thread-1' ... message was ('#RETURN', 'test') ... exception was BrokenPipeError(32, 'Broken Pipe')
问题根源
直接调用kill()会强制终止进程,此时进程内的线程还在执行Queue.get()操作——而Manager()创建的Queue依赖跨进程通信管道,进程被强制杀死后管道直接断开,线程的IO操作就会触发BrokenPipeError。
可靠实现方案
要实现进程内线程从共享Queue可靠获取数据,核心是避免强制杀死进程,改用优雅退出机制,同时在IO操作处捕获异常处理断开场景:
1. 引入退出通知事件
用multiprocessing.Event()给线程和进程传递退出信号,代替暴力kill()。
2. 捕获Queue操作的异常
在线程的get()操作处捕获BrokenPipeError和EOFError(管道断开时可能触发的异常),安全退出线程。
3. 确保进程优雅退出
进程退出前先通知线程停止,等待线程结束后再退出进程。
修改后的代码示例:
from multiprocessing import Process, Manager, Event from threading import Thread from time import sleep import multiprocessing def test_process(): def listener(q, exit_event): def worker(x, stop_event): while not stop_event.is_set(): try: # 设置超时,避免一直阻塞导致无法响应退出信号 item = x.get(True, timeout=0.5) print(f"Received: {item}") except multiprocessing.queues.Empty: continue except (BrokenPipeError, EOFError): # 管道已断开,直接退出 print("Pipe disconnected, exiting thread") break t = Thread(target=worker, args=(q, exit_event), daemon=True) t.start() # 等待线程结束(可选,根据实际场景调整) t.join() def publisher(q, exit_event): while not exit_event.is_set(): try: q.put('test') sleep(0.1) except (BrokenPipeError, EOFError): break with Manager() as mp: q = mp.Queue() exit_event = Event() a = Process(target=listener, args=(q, exit_event), daemon=True) a.start() sleep(2) b = Process(target=publisher, args=(q, exit_event), daemon=True) b.start() sleep(1) # 发送退出信号,代替kill() exit_event.set() # 等待进程优雅退出 a.join(timeout=3) b.join(timeout=3) # 如果超时还没退出,再强制杀死(作为兜底) if a.is_alive(): a.kill() if b.is_alive(): b.kill()
关键说明
- 退出事件
exit_event:让进程和线程能感知到退出信号,主动停止循环,避免暴力中断。 get()设置超时:如果不设置超时,线程会一直阻塞在get()上,无法响应退出事件,超时后会检查退出信号。- 异常捕获:处理管道断开时的
BrokenPipeError和EOFError,确保线程能安全退出,不会留下报错日志。 - 优雅退出优先:先发送退出信号并等待进程结束,超时后再强制杀死,作为最后的兜底方案。
内容的提问来源于stack exchange,提问作者Belos
相关产品推荐
相关产品推荐

