如何在Jupyter中同步采集实时传感器数据并统计特征?
实时传感器数据流处理方案需求
我通过API获取批量传感器数据(每批约200ms,采样率400Hz),已在静态数据集上用Butter低通滤波器实现了峰值统计(代码见下文)。但实际场景中需要在Jupyter Notebook里同时完成:
- 记录所有采集到的数据
- 实时检测并统计峰值(每检测到峰值执行
peaks+=1并打印当前计数)
之前尝试用两个Notebook分别处理数据记录和峰值统计:一个把数据写入sqlite3内存库,另一个读取最近5秒数据分析,但内存sqlite3无法跨Notebook共享,磁盘IO的成本又太高。
我认为最优方案是用内存中的FIFO队列+线程同时处理两项任务,想请教有没有适合API/串口实时数据的分析框架,能处理持续到来的数据并支持滑动窗口(类似data.iloc[-window_length:,:])的时序特征分析?
静态数据处理参考代码
import pandas as pd import seaborn as sns from scipy import signal # 加载静态数据 df = pd.read_excel('g 2022-09-19_17-18-56.xls') sample_rate = df.index[-1]/df.iloc[-1,0] print(sample_rate, 'Hz') # 绘制原始传感器数据 ax = sns.lineplot(data=df.iloc[:,:], x='Time (s)', y='Absolute acceleration (m/s^2)') ax.figure.set_size_inches(20,10) # Butterworth低通滤波 b, a = signal.butter(6, 0.005, 'lowpass') filtedData = signal.filtfilt(b, a, df.loc[:,'Absolute acceleration (m/s^2)']) # data为要过滤的信号 ax = sns.lineplot(x=df['Time (s)'], y=filtedData) ax.figure.set_size_inches(20,10) # 检测峰值 peaks,_ = signal.find_peaks(x=-filtedData, distance=400) ax = sns.scatterplot(x=df['Time (s)'][peaks], y=filtedData[peaks], color='purple', marker='s',) ax.figure.set_size_inches(20,10) peaks_count = len(peaks) print(f"There are {peaks_count} peaks in static data.") # 以下是待实现的实时数据逻辑 # steaming_data_=requests.get() # collecting all the data and detect # while 1: # if peaks detected: # print(peaks count)
解决方案
一、基于线程+队列的轻量实现
用Python标准库的threading和queue.Queue就能实现内存内的多任务处理,完全规避IO开销:
- 生产者线程:持续从API拉取数据,存入队列
- 消费者线程:从队列取出数据,同时完成全量数据记录、滑动窗口维护、滤波和峰值检测
示例代码
import threading import queue import time import requests import pandas as pd from scipy import signal # 配置参数 SAMPLE_RATE = 400 # Hz WINDOW_SECONDS = 5 # 滑动窗口时长 WINDOW_SIZE = SAMPLE_RATE * WINDOW_SECONDS # 窗口数据点数量 API_URL = "你的传感器API地址" peaks_count = 0 lock = threading.Lock() # 用于峰值计数的线程安全操作 # 初始化队列、全量数据存储、滑动窗口数据 data_queue = queue.Queue(maxsize=10) # 限制队列大小防止内存溢出 full_data = pd.DataFrame(columns=['Time (s)', 'Absolute acceleration (m/s^2)']) window_data = pd.DataFrame(columns=['Time (s)', 'Absolute acceleration (m/s^2)']) # 预计算Butterworth滤波器系数(和静态代码一致) b, a = signal.butter(6, 0.005, 'lowpass') # 初始化滤波器状态(用于实时滤波) zi = signal.lfilter_zi(b, a) def producer(): """生产者线程:从API拉取数据并存入队列""" while True: try: # 模拟API请求,替换为实际的API调用逻辑 response = requests.get(API_URL) batch_data = response.json() # 假设返回的是符合格式的批量数据 # 转换为DataFrame(根据实际API返回结构调整) batch_df = pd.DataFrame(batch_data, columns=['Time (s)', 'Absolute acceleration (m/s^2)']) data_queue.put(batch_df) time.sleep(0.2) # 匹配每批200ms的推送间隔 except Exception as e: print(f"API请求错误: {e}") time.sleep(1) def consumer(): """消费者线程:处理数据、记录全量、维护滑动窗口、检测峰值""" global peaks_count, full_data, window_data, zi while True: batch_df = data_queue.get() if batch_df is None: break # 1. 记录全量数据 full_data = pd.concat([full_data, batch_df], ignore_index=True) # 2. 更新滑动窗口数据 window_data = pd.concat([window_data, batch_df], ignore_index=True) # 只保留最后WINDOW_SIZE个数据点 if len(window_data) > WINDOW_SIZE: window_data = window_data.iloc[-WINDOW_SIZE:].reset_index(drop=True) # 3. 实时滤波(替换静态代码的filtfilt,用lfilter支持流式处理) filted_window, zi = signal.lfilter(b, a, window_data['Absolute acceleration (m/s^2)'], zi=zi) # 4. 检测峰值(注意distance参数需要匹配采样率) peaks, _ = signal.find_peaks(x=-filted_window, distance=400) # 统计新增峰值数(避免重复统计) new_peaks = len(peaks) - peaks_count if new_peaks > 0: with lock: peaks_count += new_peaks print(f"当前峰值计数: {peaks_count}") data_queue.task_done() # 启动线程 producer_thread = threading.Thread(target=producer, daemon=True) consumer_thread = threading.Thread(target=consumer, daemon=True) producer_thread.start() consumer_thread.start() # Jupyter中直接运行即可,无需join阻塞主线程
二、专用实时数据分析框架推荐
如果需要更灵活的流式处理能力,可以选择以下框架:
1. Streamz
专为Python设计的流式数据处理库,能无缝对接Pandas、Scipy,支持滑动窗口、聚合等操作,适配Jupyter交互环境,可快速搭建流式数据处理管道。
2. River(原Creme)
专注在线机器学习的轻量库,提供增量式时序处理工具,支持滑动窗口统计、在线滤波,内存占用极低,适合资源受限的实时场景。
3. Obspy
针对振动、地震类传感器数据的专业库,内置成熟的实时数据流处理模块,提供标准化的滤波、峰值检测算法,适合工业级传感器数据分析。
内容的提问来源于stack exchange,提问作者WilsonF
相关产品推荐
相关产品推荐

