Python 3.7多进程队列死锁排查:子进程挂起问题求助
解决多进程队列导致的子进程挂起问题
我在Docker Ubuntu环境中运行Python 3.7程序,程序包含三个通过队列连接的进程:
- 进程1(StreamCapture)向Queue 1写入视频帧数据
- 进程2(StreamProcessing)从Queue 1读取数据处理后写入Queue 2
- 进程3(StreamSender)从Queue 2读取数据并模拟发送
程序运行一段时间后,视频处理进程停止响应,exit_code为None,同时Queue 1出现过载。怀疑是同一队列的Queue.get()与Queue.put()操作阻塞导致死锁,进而引发进程挂起。程序中的config是由config.yaml生成的Pydantic模型。
原代码如下:
import multiprocessing as mp import cv2 import time from module import config class StreamCapture: def __init__(self, queue): self.queue = queue def run(self): cap = cv2.VideoCapture(config.rtsp) while cap.isOpened(): ret, frame = cap.read() if not ret: break self.queue.put(frame) time.sleep(0.01) # Simulate frame rate cap.release() class StreamProcessing: def __init__(self, input_queue, output_queue): self.input_queue = input_queue self.output_queue = output_queue def run(self): while True: frame = self.get() if frame is None: break # Example processing: convert to grayscale processed_frame = cv2.cvtColor(frame, config.colormap) self.output_queue.put(processed_frame) def get(self): ''' Function to pass some frames, because of this class is much slower then VideoCapture ''' rate = math.exp(0.3827 + 0.0449 * (self.input_queue.qsize())) rate = 100 if rate >= 100 else rate rate = 0 if rate <= 0 else rate pass_frames = int(round(rate / 100 * self.input_queue.qsize(), 0)) for _ in range(0, pass_frames): self.input_queue.get() class StreamSender: def __init__(self, queue): self.queue = queue def run(self): while True: frame = self.queue.get() if frame is None: break # Simulate sending frame to a server print("Sending frame to server") time.sleep(config.sleep_time) # Simulate network delay if __name__ == "__main__": capture_queue = mp.Queue(maxsize=10) processing_queue = mp.Queue(maxsize=10) capture = StreamCapture(capture_queue) processing = StreamProcessing(capture_queue, processing_queue) sender = StreamSender(processing_queue) capture_process=mp.Process(target=capture.run, daemon=True) processing_process=mp.Process(target=processing.run, daemon=True) sender_process=mp.Process(target=sender.run, daemon=True) capture_process.start() processing_process.start() sender_process.start() try: while True: time.sleep(1) except KeyboardInterrupt: pass finally: capture_queue.put(None) processing_queue.put(None) capture_process.join() processing_process.join() sender_process.join()
问题根源分析
- StreamProcessing.get()逻辑缺失:
- 未导入
math模块,直接运行会报错 - 跳过帧的循环未判断队列是否为空,可能永久阻塞在
input_queue.get() - 方法末尾无return语句,导致
run()中拿到的frame始终为None,处理进程会直接退出
- 未导入
- 队列操作无超时处理:
- 生产者
put()时如果队列满会永久阻塞,后续无人消费队列则引发死锁 - 消费者
get()无超时,遇到异常情况会一直挂起
- 生产者
- 终止流程有问题:
- 主线程往队列
put(None)时如果队列满,会阻塞在finally块,无法执行进程join操作 - daemon进程在主线程结束时会被强制杀死,但如果进程正阻塞在队列操作,可能无法正常释放资源
- 主线程往队列
修复方案
1. 修复StreamProcessing.get()方法
补全模块导入,添加队列空判断,确保返回有效帧:
import math # 必须添加该导入 def get(self): '''跳过部分帧,因为本类处理速度远慢于VideoCapture''' if self.input_queue.empty(): return None rate = math.exp(0.3827 + 0.0449 * (self.input_queue.qsize())) rate = min(100, max(0, rate)) pass_frames = int(round(rate / 100 * self.input_queue.qsize(), 0)) # 跳过帧时检查队列状态,避免阻塞 for _ in range(pass_frames): if self.input_queue.empty(): break try: self.input_queue.get(timeout=0.1) except mp.Queue.Empty: break # 返回当前要处理的帧 try: return self.input_queue.get(timeout=0.1) except mp.Queue.Empty: return None
2. 给队列操作添加超时与容错
- 生产者
put()时添加超时,队列满时丢弃旧帧:
# StreamCapture.run()中修改put逻辑 try: self.queue.put(frame, timeout=0.1) except mp.Queue.Full: pass
- 消费者
get()添加超时,防止意外阻塞
3. 修正进程终止逻辑
采用非阻塞方式发送终止信号,超时未退出则强制终止进程:
finally: # 非阻塞发送终止信号,避免队列满时阻塞主线程 try: capture_queue.put(None, block=False) except mp.Queue.Full: pass try: processing_queue.put(None, block=False) except mp.Queue.Full: pass # 等待进程退出,超时则强制终止 capture_process.join(timeout=5) processing_process.join(timeout=5) sender_process.join(timeout=5) if capture_process.is_alive(): capture_process.terminate() if processing_process.is_alive(): processing_process.terminate() if sender_process.is_alive(): sender_process.terminate() # 确保进程完全退出 capture_process.join() processing_process.join() sender_process.join()
4. 添加队列状态监控
在主线程中定期打印队列大小,方便排查过载问题:
try: while True: time.sleep(1) print(f"Capture queue size: {capture_queue.qsize()}, Processing queue size: {processing_queue.qsize()}") except KeyboardInterrupt: pass
完整修复后的代码
import multiprocessing as mp import cv2 import time import math from module import config class StreamCapture: def __init__(self, queue): self.queue = queue def run(self): cap = cv2.VideoCapture(config.rtsp) try: while cap.isOpened(): ret, frame = cap.read() if not ret: break # 队列满时丢弃帧,避免阻塞 try: self.queue.put(frame, timeout=0.1) except mp.Queue.Full: pass time.sleep(0.01) finally: cap.release() # 发送处理结束信号 try: self.queue.put(None, block=False) except mp.Queue.Full: pass class StreamProcessing: def __init__(self, input_queue, output_queue): self.input_queue = input_queue self.output_queue = output_queue def run(self): while True: frame = self.get() if frame is None: break processed_frame = cv2.cvtColor(frame, config.colormap) # 处理队列满时的情况 try: self.output_queue.put(processed_frame, timeout=0.1) except mp.Queue.Full: pass # 发送处理结束信号 try: self.output_queue.put(None, block=False) except mp.Queue.Full: pass def get(self): '''跳过部分帧,适配处理速度''' if self.input_queue.empty(): return None rate = math.exp(0.3827 + 0.0449 * (self.input_queue.qsize())) rate = min(100, max(0, rate)) pass_frames = int(round(rate / 100 * self.input_queue.qsize(), 0)) for _ in range(pass_frames): if self.input_queue.empty(): break try: self.input_queue.get(timeout=0.1) except mp.Queue.Empty: break try: return self.input_queue.get(timeout=0.1) except mp.Queue.Empty: return None class StreamSender: def __init__(self, queue): self.queue = queue def run(self): while True: try: frame = self.queue.get(timeout=1) if frame is None: break print("Sending frame to server") time.sleep(config.sleep_time) except mp.Queue.Empty: continue if __name__ == "__main__": capture_queue = mp.Queue(maxsize=10) processing_queue = mp.Queue(maxsize=10) capture = StreamCapture(capture_queue) processing = StreamProcessing(capture_queue, processing_queue) sender = StreamSender(processing_queue) capture_process = mp.Process(target=capture.run) processing_process = mp.Process(target=processing.run) sender_process = mp.Process(target=sender.run) capture_process.start() processing_process.start() sender_process.start() try: while True: time.sleep(1) print(f"Capture queue size: {capture_queue.qsize()}, Processing queue size: {processing_queue.qsize()}") except KeyboardInterrupt: print("Received interrupt, stopping processes...") finally: try: capture_queue.put(None, block=False) except mp.Queue.Full: pass try: processing_queue.put(None, block=False) except mp.Queue.Full: pass capture_process.join(timeout=5) processing_process.join(timeout=5) sender_process.join(timeout=5) if capture_process.is_alive(): capture_process.terminate() if processing_process.is_alive(): processing_process.terminate() if sender_process.is_alive(): sender_process.terminate() capture_process.join() processing_process.join() sender_process.join()
内容的提问来源于stack exchange,提问作者Nikita Belov
相关产品推荐
相关产品推荐

