Python与Arduino Uno高数据速率并行处理方案求助
EMG数据采集项目多进程并行方案解决思路
核心问题分析
- 线程方案受GIL限制无法实现真正并行,0.03秒的处理耗时会阻塞采集线程,导致样本丢失
- 多进程方案中PySerial对象无法被pickle序列化,直接跨进程传递会触发
ValueError: ctypes objects containing pointers cannot be pickled错误
可行解决方案:进程间通信+独立串口初始化
不要在主进程创建串口后传递给子进程,而是让采集进程独立初始化串口,通过进程安全的队列(multiprocessing.Queue)传递采集数据,处理进程从队列取数完成预处理和分类。完全规避串口对象的序列化问题,同时保证数据不丢失。
具体实现步骤
- 创建进程安全队列作为数据缓冲区,限制最大容量避免内存溢出
- 采集进程:独立初始化串口,按250Hz频率读取数据并写入队列
- 处理进程:从队列批量取数,积累足够样本后执行预处理和分类,剩余样本留待下一次循环处理
代码示例
import multiprocessing import serial import time import numpy as np # 采集进程函数 def emg_collector(queue, port='COM3', baud_rate=9600): # 采集进程内独立初始化串口 ser = serial.Serial(port, baud_rate) ser.flushInput() try: while True: # 严格控制250Hz采集频率(每0.004秒一次) cycle_start = time.perf_counter() if ser.in_waiting > 0: raw_data = ser.readline().decode('utf-8').strip() try: emg_sample = float(raw_data) # 队列满时阻塞,避免丢数据 queue.put(emg_sample) except ValueError: # 跳过无效数据 pass # 补全剩余时间,保证采集频率稳定 elapsed = time.perf_counter() - cycle_start if elapsed < 0.004: time.sleep(0.004 - elapsed) except KeyboardInterrupt: ser.close() print("采集进程终止") # 处理进程函数 def emg_processor(queue): # 单次处理对应0.03秒数据(250Hz下共75个样本) process_window = 75 sample_buffer = [] try: while True: # 批量读取队列中所有可用数据 while not queue.empty(): sample_buffer.append(queue.get()) # 缓冲区样本量足够时执行处理 if len(sample_buffer) >= process_window: # 截取完整窗口数据,剩余样本留到下一次 current_batch = sample_buffer[:process_window] sample_buffer = sample_buffer[process_window:] # 替换为你的预处理逻辑(滤波、归一化等) processed_data = np.array(current_batch) * 0.5 # 替换为你的分类逻辑 intent_label = np.argmax([np.mean(processed_data), np.std(processed_data)]) print(f"识别意图: {intent_label}, 处理耗时: {time.perf_counter()-start:.4f}秒") # 降低CPU占用,给采集进程留运行时间 time.sleep(0.001) except KeyboardInterrupt: print("处理进程终止") if __name__ == "__main__": # 创建队列,设置最大容量防止内存溢出 data_queue = multiprocessing.Queue(maxsize=1000) # 启动采集进程 collector_proc = multiprocessing.Process(target=emg_collector, args=(data_queue,)) collector_proc.start() # 启动处理进程 processor_proc = multiprocessing.Process(target=emg_processor, args=(data_queue,)) processor_proc.start() # 等待用户中断程序 try: while True: time.sleep(1) except KeyboardInterrupt: collector_proc.terminate() processor_proc.terminate() collector_proc.join() processor_proc.join() print("所有进程已终止")
关键注意事项
- 串口独立初始化:每个进程自行创建串口对象,彻底避免跨进程传递序列化问题
- 队列阻塞机制:
queue.put()默认阻塞,当队列满时采集进程会暂停,直到处理进程取走数据,保证样本不丢失 - 样本完整性:处理时按固定窗口截取数据,剩余样本保留在缓冲区,确保无样本遗漏
- 资源释放:添加异常捕获逻辑,程序终止时正常关闭串口并回收进程资源
额外优化建议
- 用
time.perf_counter()替代time.time(),提升时间控制精度 - 预处理和分类逻辑采用NumPy向量化操作,进一步压缩处理耗时
- 可给队列添加超时取数机制,避免处理进程长时间空转
内容的提问来源于stack exchange,提问作者Indrayudd Roy Chowdhury
相关产品推荐
相关产品推荐

