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

进程中启动的线程过早终止问题排查求助

问题排查:多进程下RGBD帧接收/保存线程异常终止

我正在搭建一套由3台连接RGBD相机的Raspberry Pi组成的同步录制系统。原本每台RPi发送的帧由recv_data线程接收,接收的帧再由独立的save线程保存,共创建6个线程且运行正常。但改为将每组接收/保存线程放入独立进程后,线程在while True的首次循环后就停止运行,附上代码,恳请帮忙排查问题原因。

def recv_data(socket, display_title, q_color, q_depth, q_timestamps):
    # i is the frames counter
    i = 1

    print('Receiving frames from ' + display_title)

    while True:
        # receive color_image, depth_image and timestamp
        color_msg, depth_msg, timestamp = socket.recv_multipart()

        # decompress qoi
        color_msg = qoi.decode(color_msg)

        # convert and reshape images and decode timestamp to string
        color_frame = np.frombuffer(color_msg, 'uint8')
        depth_frame = np.frombuffer(depth_msg, 'uint16')
        timestamp = timestamp.decode()

        print(display_title)

        q_color.put(color_frame)
        q_depth.put(depth_frame)
        q_timestamps.put(timestamp)

        # increment index
        i = i + 1


# thread to save the frames from color frames and depth frames global lists
def save(display_title, q_color, q_depth, q_timestamps):
    i = 1

    while True:
        if not q_color.empty():

            # save images in numpy array
            path = frame_path + '/' + display_title + '/'
            save_frame(q_color.get(), 'color', path, i)
            save_frame(q_depth.get(), 'depth', path, i)

            # save timestamp in csv file
            csv1 = frame_path + '/' + display_title.replace(" ", "") + '.csv'
            data = open(csv1, 'a', newline='')
            writer_data = csv.writer(data)
            save_timestamp(writer_data, q_timestamps.get(), i)

            # increment counter
            i = i + 1


# Acquisition function
def start_acquisition():
    # connect server to client
    socket1 = server_connect('tcp://192.168.1.150:5551', 'RPi1')
    socket2 = server_connect('tcp://192.168.1.144:5552', 'RPi2')
    socket3 = server_connect('tcp://192.168.1.169:5553', 'RPi3')

    # initiate csv file for timestamps
    file_name = frame_path + '/RPi1.csv'
    open(file_name, 'w').close()
    file_name = frame_path + '/RPi2.csv'
    open(file_name, 'w').close()
    file_name = frame_path + '/RPi3.csv'
    open(file_name, 'w').close()

    p1 = Process(target=start_acquisition_threads, args=(socket1, "RPi 1"))
    p2 = Process(target=start_acquisition_threads, args=(socket2, "RPi 2"))
    p3 = Process(target=start_acquisition_threads, args=(socket3, "RPi 3"))

    p1.start()
    p2.start()
    p3.start()

    p1.join()
    p2.join()
    p3.join()


def start_acquisition_threads(socket, rpi):
    # define queues to manage frames and timestamps
    q_color = Queue()
    q_depth = Queue()
    q_timestamps = Queue()

    # create and start threads
    th_recv_data = Thread(target=recv_data, args=(socket, rpi, q_color, q_depth, q_timestamps))
    th_save_data = Thread(target=save, args=(rpi,q_color, q_depth, q_timestamps))

    th_recv_data.start()
    th_save_data.start()

核心问题分析

  1. 子进程过早退出导致线程被终止
    start_acquisition_threads函数启动两个线程后没有做阻塞等待,函数执行完毕后对应的子进程直接退出,进程内所有线程会被强制终止,这就是线程只跑了首次循环就停止的直接原因。

  2. 跨进程传递Socket对象无效
    主进程中创建的Socket对象无法直接跨进程序列化传递,子进程拿到的Socket处于无效状态,会导致recv_multipart调用失败(未捕获异常直接终止线程)。

修复方案

方案1:让子进程保持存活,等待线程运行

修改start_acquisition_threads,启动线程后加入阻塞等待:

def start_acquisition_threads(socket, rpi):
    q_color = Queue()
    q_depth = Queue()
    q_timestamps = Queue()

    th_recv_data = Thread(target=recv_data, args=(socket, rpi, q_color, q_depth, q_timestamps))
    th_save_data = Thread(target=save, args=(rpi,q_color, q_depth, q_timestamps))

    th_recv_data.start()
    th_save_data.start()
    
    # 阻塞等待线程结束(因线程是while True,会一直保持进程存活)
    th_recv_data.join()
    th_save_data.join()

方案2:子进程内创建Socket,避免跨进程传递

把Socket创建逻辑移到子进程内部:

def start_acquisition_threads(addr, rpi):
    # 子进程内独立创建Socket
    socket = server_connect(addr, rpi)
    
    q_color = Queue()
    q_depth = Queue()
    q_timestamps = Queue()

    th_recv_data = Thread(target=recv_data, args=(socket, rpi, q_color, q_depth, q_timestamps))
    th_save_data = Thread(target=save, args=(rpi,q_color, q_depth, q_timestamps))

    th_recv_data.start()
    th_save_data.start()
    
    th_recv_data.join()
    th_save_data.join()

# 修改start_acquisition中的进程创建代码
def start_acquisition():
    # 初始化csv文件逻辑不变
    
    p1 = Process(target=start_acquisition_threads, args=('tcp://192.168.1.150:5551', "RPi 1"))
    p2 = Process(target=start_acquisition_threads, args=('tcp://192.168.1.144:5552', "RPi 2"))
    p3 = Process(target=start_acquisition_threads, args=('tcp://192.168.1.169:5553', "RPi 3"))

    p1.start()
    p2.start()
    p3.start()

    p1.join()
    p2.join()
    p3.join()

补充优化建议

  • 给recv_data和save线程添加异常捕获,避免单次帧处理失败导致线程退出:
    def recv_data(socket, display_title, q_color, q_depth, q_timestamps):
        i = 1
        print('Receiving frames from ' + display_title)
        while True:
            try:
                # 原有接收处理逻辑
            except Exception as e:
                print(f"{display_title}接收线程异常: {e}")
                time.sleep(0.1)
    
  • save函数中频繁打开关闭CSV文件影响性能,建议在函数开头打开文件并保持,或使用锁机制避免并发写入问题。

内容的提问来源于stack exchange,提问作者BoolMimi

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.26 23:30:42