PySpark中无时间戳,基于frame_number的滚动平均值计算问询
基于PySpark和frame_number列计算滚动平均值
没问题!既然你的数据里有严格递增的frame_number,那它完全可以替代时间戳来实现滚动平均——毕竟滚动计算的核心就是依赖有序的序列,而frame_number刚好满足这个要求。下面我给你一步步演示具体怎么做:
1. 准备工作:初始化SparkSession并加载数据
首先我们需要导入必要的模块,把你的示例数据转换成Spark DataFrame:
from pyspark.sql import SparkSession from pyspark.sql.window import Window from pyspark.sql import functions as F # 初始化SparkSession spark = SparkSession.builder.appName("RollingAvgWithFrameNumber").getOrCreate() # 补全后的示例数据 d = [ {'session_id': 1, 'frame_number': 1, 'rtd': 11.0, 'rtd2': 11.0}, {'session_id': 1, 'frame_number': 2, 'rtd': 12.0, 'rtd2': 6.0}, {'session_id': 1, 'frame_number': 3, 'rtd': 4.0, 'rtd2': 233.0}, {'session_id': 1, 'frame_number': 4, 'rtd': 110.0, 'rtd2': 15.0}, {'session_id': 2, 'frame_number': 1, 'rtd': 8.0, 'rtd2': 9.0}, {'session_id': 2, 'frame_number': 2, 'rtd': 15.0, 'rtd2': 22.0} ] # 创建DataFrame df = spark.createDataFrame(d) df.show()
2. 定义窗口规范
滚动平均的关键是窗口范围的定义,这里我们需要按session_id分区(保证每个会话的计算独立),然后按frame_number排序(确保序列是递增的)。根据你的需求,有两种常见的窗口定义方式:
方式一:基于固定行数的滚动平均(比如最近N个帧)
如果你想计算当前帧加上前2个帧的滚动平均(总共3个样本),可以用rowsBetween来指定行的范围:
# 窗口规范:按session分组,按frame_number排序,取当前行及前2行 window_spec_rows = Window.partitionBy("session_id").orderBy("frame_number").rowsBetween(-2, 0) # 计算rtd和rtd2的滚动平均值(保留2位小数) df_rolling_rows = df.withColumn("rtd_rolling_avg", F.round(F.avg("rtd").over(window_spec_rows), 2)) \ .withColumn("rtd2_rolling_avg", F.round(F.avg("rtd2").over(window_spec_rows), 2)) df_rolling_rows.show()
rowsBetween(-2, 0)的意思是:窗口包含从当前行往前数2行(-2)到当前行(0)的所有数据。- 用
F.round()是为了让结果更整洁,你可以根据需求调整小数位数。
方式二:基于frame_number数值范围的滚动平均
如果你的frame_number可能存在缺失(比如不是连续的整数),或者你想基于frame_number的差值来定义窗口(比如只包含与当前帧的frame_number差≤2的所有帧),可以用rangeBetween:
# 窗口规范:按session分组,按frame_number排序,取frame_number在[当前帧-2, 当前帧]之间的行 window_spec_range = Window.partitionBy("session_id").orderBy("frame_number").rangeBetween(-2, 0) # 计算滚动平均值(保留2位小数) df_rolling_range = df.withColumn("rtd_rolling_avg_range", F.round(F.avg("rtd").over(window_spec_range), 2)) \ .withColumn("rtd2_rolling_avg_range", F.round(F.avg("rtd2").over(window_spec_range), 2)) df_rolling_range.show()
这种方式的优势是:即使frame_number不连续(比如存在1、3、4这样的情况),它也会严格按照frame_number的数值范围来筛选数据,而不是单纯取前N行。
3. 注意事项
- 必须按
session_id分区:否则不同会话的帧会被混在一起计算,结果就错了。 - 必须按
frame_number排序:窗口函数依赖有序的序列才能正确定义滚动范围。 - 灵活调整窗口范围:如果需要计算最近5个帧的平均,把
rowsBetween(-4, 0)即可(因为0是当前行,-4代表往前4行,总共5行)。 - 支持其他聚合函数:除了
avg(),你还可以用sum()、max()、min()等函数,实现滚动求和、滚动最大/最小值等。
内容的提问来源于stack exchange,提问作者Prince Bhatti
相关产品推荐
相关产品推荐

