Python多进程架构咨询:传感器数据采集与批量读取并行方案
问题解答
1. 多进程内存共享问题
Python multiprocessing 模块中,默认进程间不共享内存——每个子进程会复制父进程的内存空间,后续各自独立修改,互不影响。如果直接用普通列表作为Buffer,进程1追加的数据进程2完全无法读取。
要实现进程间数据共享,常用两种方案:
- 用
multiprocessing.Queue:进程安全的消息队列,自带同步机制,完美适配生产者-消费者场景(你的需求正好是进程1生产数据、进程2消费数据)。 - 用
multiprocessing.Manager创建共享列表:通过管理进程代理实现数据共享,但需要手动配合锁保证操作原子性,比如取前140条时避免被进程1同时修改。
2. 最优并行架构实现
推荐采用生产者-消费者模式,基于Queue实现,无需手动加锁,代码简洁可靠。以下是完整可运行示例:
代码实现
import multiprocessing import serial import time # 进程1:串口数据采集(生产者) def serial_collector(queue, port='COM3', baudrate=9600): try: ser = serial.Serial(port, baudrate, timeout=1) print(f"已连接串口 {port}") while True: line = ser.readline().decode('utf-8').strip() if line: # 确保读取到有效数据 queue.put(line) time.sleep(0.01) # 按需调整采集间隔 except serial.SerialException as e: print(f"串口连接失败或断开: {e}") finally: if 'ser' in locals(): ser.close() # 进程2:数据处理(消费者) def data_processor(queue): buffer = [] while True: # 从队列中获取数据,填充本地buffer while not queue.empty(): buffer.append(queue.get()) # 当buffer积累到140条时执行后续任务 if len(buffer) >= 140: process_data = buffer[:140] # 替换为你的实际处理逻辑 print(f"处理140条数据,第一条: {process_data[0]},最后一条: {process_data[-1]}") # 保留buffer中未处理的剩余数据 buffer = buffer[140:] time.sleep(0.1) # 避免空轮询占用过多CPU if __name__ == '__main__': # 创建带容量限制的进程安全队列,防止内存溢出 data_queue = multiprocessing.Queue(maxsize=1000) # 创建并启动进程 collector_process = multiprocessing.Process(target=serial_collector, args=(data_queue,)) processor_process = multiprocessing.Process(target=data_processor, args=(data_queue,)) collector_process.start() processor_process.start() # 等待进程运行,手动终止时清理 try: collector_process.join() processor_process.join() except KeyboardInterrupt: print("程序手动终止") collector_process.terminate() processor_process.terminate()
架构说明
- 生产者进程:负责串口数据读取,将有效数据存入
Queue,串口断开时自动捕获异常并释放资源。 - 消费者进程:持续从队列取数据填充本地buffer,达到140条阈值时取出处理,剩余数据继续留存。
- Queue优势:自带进程同步,无需手动加锁;设置
maxsize可限制队列长度,避免采集速度远超处理速度导致内存爆炸。
备选方案(共享列表+锁)
如果坚持用共享列表,可通过Manager创建并配合锁保证安全,示例片段:
from multiprocessing import Manager, Lock def serial_collector(shared_list, lock): while True: line = ser.readline().decode().strip() if line: with lock: shared_list.append(line) def data_processor(shared_list, lock): while True: with lock: if len(shared_list) >= 140: process_data = shared_list[:140] del shared_list[:140] # 执行process_data的后续处理逻辑
这种方式需手动管理锁,易出现死锁或数据不一致,优先推荐Queue方案。
3. 多核心并行说明
multiprocessing创建的子进程,操作系统会自动调度到不同CPU核心运行——多进程绕过了Python的GIL(全局解释器锁)限制,真正实现并行执行,无需额外配置。
内容的提问来源于stack exchange,提问作者CodeChefLearner
相关产品推荐
相关产品推荐

