如何在PySpark中实现基于月份而非日期的rangeBetween窗口计算
在PySpark中实现基于月份范围的滑动窗口均值计算
由于月份天数不固定,无法直接用datediff配合rangeBetween做数值范围的窗口限制,这里提供两种实用方案:
方案一:按自然月对齐的窗口计算
如果业务需要包含当前行所在月份及往前3个完整自然月的范围(例如当前日期为2024-05-15,范围则是2024-02-01至2024-05-15),可通过日期截断+月份差过滤实现:
from pyspark.sql import functions as F from pyspark.sql.window import Window # 假设数据集为df,包含date_col(日期类型)、some_value及分组键group_key df = df.withColumn( "month_start", F.date_trunc("month", F.col("date_col")) # 将日期截断至当月第一天 ) # 定义窗口:按分组键分区,按日期排序,窗口覆盖当前行及之前所有行 window_spec = Window.partitionBy("group_key") .orderBy(F.col("date_col").cast("long")) .rangeBetween(Window.unboundedPreceding, 0) # 过滤出近3个月内的数据并计算均值 result_df = df.withColumn( "avg_3months", F.avg( F.when( F.months_between(F.col("month_start"), F.date_trunc("month", F.col("date_col"))).between(-3, 0), F.col("some_value") ) ).over(window_spec) )
方案二:精确日期范围的窗口计算
如果需要当前行日期往前推3个月的精确区间(例如当前日期为2024-05-15,范围则是2024-02-15至2024-05-15),可通过动态计算起始日期+窗口过滤实现:
from pyspark.sql import functions as F from pyspark.sql.window import Window # 计算每行对应的3个月前的起始日期 df = df.withColumn( "start_date", F.add_months(F.col("date_col"), -3) ) # 定义窗口:按分组键分区,按日期排序,覆盖当前行及之前所有行 window_spec = Window.partitionBy("group_key") .orderBy(F.col("date_col").cast("long")) .rowsBetween(Window.unboundedPreceding, Window.currentRow) # 收集窗口内数据并过滤出符合日期范围的数值,再计算均值 result_df = df.withColumn( "window_data", F.collect_list(F.struct("date_col", "some_value")).over(window_spec) ).withColumn( "valid_data", F.expr("filter(window_data, x -> x.date_col >= start_date)") ).withColumn( "avg_3months", F.avg(F.col("valid_data.some_value")) ).drop("window_data", "valid_data", "start_date")
注意事项
- 若数据量较大,方案二的
collect_list可能引发性能问题,优先选择方案一(业务逻辑允许按自然月对齐时)。 - 确保
date_col为PySpark的DateType或TimestampType,字符串类型需先转换:F.to_date(F.col("date_str"), "yyyy-MM-dd")
内容的提问来源于stack exchange,提问作者Tinkerbell
相关产品推荐
相关产品推荐

