Python实现流式数据移动平均及相邻移动平均差值统计方法
流式数据移动平均及差值计数实现方案
你原本用pandas实现的离线逻辑可以完全对齐到流式场景,只需要维护3个持久化的状态变量,逐条处理流入的数据即可:
所需维护的状态变量
- 固定长度为5的滑动窗口容器,用来存储最近到达的5个S1字段值
- 上一次计算得到的移动平均结果,初始值设为空
- 差值大于4的事件计数器,初始值设为0
单条数据处理流程
每流入一条新的S1数据时,按以下顺序执行操作:
- 将新的S1值加入滑动窗口容器
- 若窗口内元素数量不足5,直接结束当前处理流程,不做计算
- 若窗口内元素数量等于5:
- 计算当前窗口的平均值
current_avg = sum(窗口内所有元素) / 5 - 如果上一次移动平均结果不为空,计算两者差值
diff = current_avg - 上一次移动平均结果,如果差值的绝对值大于4(若只需要判断正差值可以去掉绝对值),计数器加1 - 将上一次移动平均结果更新为当前计算的
current_avg - 弹出窗口中最早存入的元素,保持窗口长度为5,等待下一条数据流入
- 计算当前窗口的平均值
Python 原生实现示例
from collections import deque # 初始化状态 window = deque(maxlen=5) # 定长队列,自动弹出最早元素,无需手动处理 last_moving_avg = None count = 0 def process(new_s1): global last_moving_avg, count window.append(new_s1) if len(window) < 5: return current_avg = sum(window) / 5 if last_moving_avg is not None: diff = current_avg - last_moving_avg if abs(diff) > 4: count += 1 last_moving_avg = current_avg # 流式调用方式:每收到一条数据调用一次process函数即可 # 例:收到S1值为12时,执行 process(12)
如果你用的是Flink、Spark Streaming等专业流式计算框架,直接调用框架内置的滑动窗口API设置窗口大小为5,再通过框架的状态功能存储上一次窗口均值做差值判断即可,逻辑和上述实现完全一致。
内容的提问来源于stack exchange,提问作者Edward
相关产品推荐
相关产品推荐

