Pyspark按企业分组计算12个月移动平均值代码出错如何解决
问题原因
- 你当前使用的
rowsBetween(-11, 0)是物理行偏移规则,仅会取当前行往前数11行+当前行共12行数据计算平均,完全不识别时间范围:- 你的样例数据中同一企业同一日期存在多条重复记录,12行数据实际覆盖的时间远小于12个月
- 即使没有重复日期,如果存在月份断档,12行数据的时间跨度也会超过12个月,不符合需求
- 未显式将
CALENDAR_DATE转为日期类型,直接按字符串排序可能出现日期顺序错误的问题
正确实现方案
场景1:保留所有原始行,每行对应计算其所属日期往前12个月的所有数据平均值
from pyspark.sql.functions import * from pyspark.sql.window import Window # 1. 先把日期字段转成标准日期格式 df = df.withColumn('CALENDAR_DATE', to_date(col('CALENDAR_DATE'))) # 2. 计算每个日期对应的月份数值(用于时间范围偏移计算) df = df.withColumn('month_offset', months_between(col('CALENDAR_DATE'), lit('1970-01-01')).cast('int')) # 3. 定义基于时间范围的窗口:按企业分组,按月份偏移排序,取往前11个月到当月的所有数据 w = Window.partitionBy('COMPANY').orderBy('month_offset').rangeBetween(-11, 0) # 4. 计算滚动平均 df = df.withColumn('ROLLING_AVERAGE', round(avg('VALUE').over(w), 1))
场景2:先按企业+月份聚合,计算每个月的单月均值/总和后,再算12个月滚动平均
如果你的需求是每个企业每月仅输出一条滚动平均结果,可以先做聚合再计算:
from pyspark.sql.functions import * from pyspark.sql.window import Window # 1. 转日期+聚合单月数据(示例取单月均值,也可以替换为sum等聚合逻辑) df_monthly = df.withColumn('CALENDAR_DATE', to_date(col('CALENDAR_DATE'))) \ .groupBy('COMPANY', 'CALENDAR_DATE') \ .agg(avg('VALUE').alias('MONTHLY_VALUE')) # 2. 计算月份偏移 df_monthly = df_monthly.withColumn('month_offset', months_between(col('CALENDAR_DATE'), lit('1970-01-01')).cast('int')) # 3. 定义窗口计算滚动平均 w = Window.partitionBy('COMPANY').orderBy('month_offset').rangeBetween(-11, 0) df_monthly = df_monthly.withColumn('ROLLING_AVERAGE', round(avg('MONTHLY_VALUE').over(w), 1))
内容的提问来源于stack exchange,提问作者cici
相关产品推荐
相关产品推荐

