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
相关产品推荐
相关产品推荐

