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

Python多进程类方法间共享numpy数组及标志位方案问询

跨进程传递标志位与Numpy数组实现方案

核心选型

多进程内存空间完全隔离,直接修改类属性无法在两个进程间同步状态,使用multiprocessing内置的两个原语即可满足需求,无需引入第三方依赖:

  • mp.Queue:负责跨进程传递numpy数组,自带阻塞机制,接收方未取走数据时发送方会自动阻塞,天然匹配"接收方未就绪就阻塞传输"的需求,无需手写标志位判断
  • mp.Event:负责跨进程同步布尔状态标志,自带阻塞等待方法,替代原有忙等轮询逻辑,等待时不占用CPU资源

最小改动实现

整体改动仅涉及通信相关的逻辑,完全不触碰传感器采集、数据处理的核心业务代码,其中驱动侧改动量不超过10行,也可以通过继承/猴子补丁实现零侵入原驱动代码。

1. 驱动侧适配

仅需替换原有本地标志位、新增数据入队逻辑,原有采集流程完全保留:

如果完全不能修改驱动源码,可通过子类继承重写初始化方法、数据输出节点的方式注入通信对象,不需要改动原有驱动文件。

import threading
import time

class DataCaptureThread():
    def __init__(self,duration,parameters, data_queue, can_send_event, can_recv_event):
        self.duration = 0
        self.parameters = parameters
        self.capture_flag = False
        # 替换原本地标志位为跨进程共享Event,原有is_set()判断逻辑和布尔值用法完全一致
        self.outgoing_data_open_flag = can_send_event
        # 注入跨进程通信对象
        self.data_queue = data_queue
        self.can_recv_event = can_recv_event
        self.thread = threading.Thread(target=self.read_data, args=(duration,))
        
    def read_data(self, duration):
        if self.capture_flag:
            while True:              
                # 原有采集逻辑完全不动
                self.data_array = ReadSocket()
                # 原忙等逻辑保留,加短休眠避免空转占满CPU
                while not self.outgoing_data_open_flag.is_set():
                    time.sleep(0.001)
                # 数据入队,队列满时自动阻塞,无需自行判断data_full状态
                self.data_queue.put(self.data_array)
                # 通知处理方数据已就绪
                self.can_recv_event.set()
                # 重置发送允许标志,等待处理方完成当前任务后再放行下一次发送
                self.outgoing_data_open_flag.clear()

# SensorInterface类原有逻辑完全不需要改,仅需在实例化后给内部capture_stream注入通信对象即可
class SensorInterface():
    def __init__(self,parameters):        
        self.parameters = parameters
        
    def data_stream_start(self):
        self.capture_stream = DataCaptureThread()
        
    def collect_data(self, duration):
        self.capture_stream.duration = duration
        self.capture_stream.capture_flag = True
        self.capture_stream.thread.start()

2. DataHandler改造

所有改动均在数据处理模块内完成,不依赖驱动侧修改:

import numpy as np
class DataHandler():
    def __init__(self, data_queue, can_send_event, can_recv_event):
        # 替换原本地类属性标志位为跨进程共享对象
        self.data_queue = data_queue
        self.can_send_event = can_send_event
        self.can_recv_event = can_recv_event
        self.data_array = np.array([])
        
    def run_processing(self):
        while True:
            # 阻塞等待采集端的数据就绪通知,无需空转轮询
            self.can_recv_event.wait()
            # 从队列取数据,取操作会自动消费队列内容
            self.data_array = self.data_queue.get()
            # 通知采集端可以开始准备下一组数据,此时当前进程会并行处理当前数据,实现采集/处理流水线
            self.can_send_event.set()
            # 原有处理逻辑完全不动
            DataProcessing(self.data_array)
            # 重置接收标志
            self.can_recv_event.clear()

3. 主函数逻辑调整

在主进程初始化共享通信对象,分别注入两个业务类即可:

if __name__ == '__main__':
    import multiprocessing as mp
    # 初始化通信原语,队列长度设为1即可,刚好缓存1组待处理数据
    data_queue = mp.Queue(maxsize=1)
    can_send_event = mp.Event()
    can_recv_event = mp.Event()
    # 初始状态允许采集端发送第一组数据
    can_send_event.set()

    # 初始化业务实例
    sensor = SensorInterface(params1)
    sensor.data_stream_start()
    # 直接给驱动内部的采集线程注入通信对象,不需要改SensorInterface类代码
    sensor.capture_stream.data_queue = data_queue
    sensor.capture_stream.outgoing_data_open_flag = can_send_event
    sensor.capture_stream.can_recv_event = can_recv_event

    data_processor = DataHandler(data_queue, can_send_event, can_recv_event)
    
    duration = 3.0   
    # 启动两个并行进程
    collect_proc = mp.Process(target=sensor.collect_data, args=(duration,)) 
    data_proc = mp.Process(target=data_processor.run_processing, args=())    
    collect_proc.start()
    data_proc.start()
    collect_proc.join()
    data_proc.join()

注意事项

  • numpy数组通过mp.Queue传递时会自动完成序列化/反序列化,无需手动做格式转换,传输效率可满足常规传感器高吞吐量采集需求
  • 禁止使用while not flag: continue的空转轮询逻辑,会占满单个CPU核心,优先使用Queue、Event自带的阻塞方法,等待状态下进程处于休眠态,几乎不消耗CPU资源
  • 队列maxsize=1的配置天然实现了"处理方忙碌时阻塞采集端发送"的逻辑,不需要额外维护data_full标志位,队列内置的满阻塞机制已经覆盖该需求
  • 如果对传输延迟有极高要求,可替换为multiprocessing.shared_memory实现numpy数组零拷贝传递,但常规场景下Queue的性能已经足够,且实现更简单不易出错

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.27 20:09:17