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

Python与Arduino Uno高数据速率并行处理方案求助

EMG数据采集项目多进程并行方案解决思路

核心问题分析

  • 线程方案受GIL限制无法实现真正并行,0.03秒的处理耗时会阻塞采集线程,导致样本丢失
  • 多进程方案中PySerial对象无法被pickle序列化,直接跨进程传递会触发ValueError: ctypes objects containing pointers cannot be pickled错误

可行解决方案:进程间通信+独立串口初始化

不要在主进程创建串口后传递给子进程,而是让采集进程独立初始化串口,通过进程安全的队列(multiprocessing.Queue)传递采集数据,处理进程从队列取数完成预处理和分类。完全规避串口对象的序列化问题,同时保证数据不丢失。

具体实现步骤

  1. 创建进程安全队列作为数据缓冲区,限制最大容量避免内存溢出
  2. 采集进程:独立初始化串口,按250Hz频率读取数据并写入队列
  3. 处理进程:从队列批量取数,积累足够样本后执行预处理和分类,剩余样本留待下一次循环处理

代码示例

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.31 03:51:36