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

如何在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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.18 12:20:40