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

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()

问题根源分析

  1. StreamProcessing.get()逻辑缺失:
    • 未导入math模块,直接运行会报错
    • 跳过帧的循环未判断队列是否为空,可能永久阻塞在input_queue.get()
    • 方法末尾无return语句,导致run()中拿到的frame始终为None,处理进程会直接退出
  2. 队列操作无超时处理:
    • 生产者put()时如果队列满会永久阻塞,后续无人消费队列则引发死锁
    • 消费者get()无超时,遇到异常情况会一直挂起
  3. 终止流程有问题:
    • 主线程往队列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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.20 02:46:06