PySpark统计窗口函数计算移动统计量时持续返回NULL值的问题排查
PySpark统计窗口函数计算移动统计量时持续返回NULL值的问题排查
嘿,我来帮你揪出这个问题的根源!从你的代码和输出结果来看,所有统计量全返回NULL的核心原因是窗口函数的范围设置完全错误,导致每个窗口里根本没有可计算的数据,统计函数自然输出NULL。
问题到底出在哪?
你写的窗口定义是这样的:
window = Window.orderBy("new_datetime").rowsBetween(5, Window.currentRow)
这里你完全搞反了PySpark中rowsBetween(start, end)的整数偏移规则:负数代表当前行之前的行,正数代表当前行之后的行。你设置的start=5,意思是"当前行之后的第5行",但你的数据是按new_datetime升序排列的(从早到晚),处理每一行时,后面的"未来数据"还没被纳入窗口范围,所以每个窗口里都没有任何数据,统计函数当然算不出结果。
另外,你想要的是7天移动窗口(当前行+前6行,共7个样本),正确的偏移应该是从当前行之前的第6行到当前行。
直接可用的修正方案
把窗口定义改成下面这样,就能得到你想要的7天移动窗口:
from pyspark.sql.window import Window from pyspark.sql.functions import skewness, stddev, avg # 正确的7天移动窗口:当前行 + 前6行,共7个样本 window = Window.orderBy("new_datetime").rowsBetween(-6, Window.currentRow) # 同时计算移动标准差、偏度和均值 result = result.withColumn('moving_std_close', stddev('closing_price').over(window)) \ .withColumn('skew_close', skewness('closing_price').over(window)) \ .withColumn('moving_avg_close', avg('closing_price').over(window)) # 查看结果 result.select('new_datetime', 'closing_price', 'moving_std_close', 'skew_close', 'moving_avg_close').show()
额外排查点(防止后续踩坑)
如果修正后还是有异常,可以检查这几个细节:
- 字段类型是否正确:确认
closing_price是数值类型(比如DoubleType),如果是字符串类型,统计函数会直接返回NULL,需要先通过cast("double")转换。 - 时间字段排序是否靠谱:确保
new_datetime是TimestampType,而不是字符串类型的时间(字符串排序会出现"2024-01-10"排在"2024-01-2"前面的错误)。 - 样本量是否达标:前6行的窗口样本数不足7个(比如第1行只有1个样本,第2行2个...第6行6个),偏度这类统计量需要至少3个样本才能计算,前几行可能还是会返回NULL,但从第7行开始就会有正常结果,这是合理的。
预期效果
修正窗口后,你会看到从第7行(索引从0开始的话是第6行)开始,moving_std_close和skew_close会返回正常数值,前6行因样本量不足出现的少量NULL是正常现象。
备注:内容来源于stack exchange,提问作者Mig Rivera Cueva
相关产品推荐
相关产品推荐

