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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.17 13:05:37