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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.25 03:24:30