PySpark优雅实现DataFrame多组合时间范围金额均值计算
优雅实现PySpark按类型统计多时间窗口金额均值
步骤1:数据类型转换
首先需要将字符串类型的金额和日期转换为数值、日期类型,方便后续计算:
from pyspark.sql import functions as F from pyspark.sql import Row # 原始数据 df = spark.createDataFrame([ Row(ttype='C', amt='12.99', dt='2024/01/01'), Row(ttype='D', amt='21.99', dt='2024/02/15'), Row(ttype='C', amt='16.99', dt='2024/01/21'), ]) # 转换数据类型 df_clean = df.withColumn("amt", F.col("amt").cast("double")) \ .withColumn("dt", F.to_date("dt", "yyyy/MM/dd"))
步骤2:批量计算固定时间窗口的均值(基于基准日期)
如果需求是统计每个类型在基准日期(如数据最大日期、当前日期)过去N天内的金额均值,可以通过定义时间窗口列表+列表推导式批量生成聚合逻辑,避免重复代码:
# 定义需要计算的时间窗口(天数) window_days = [30, 60, 90] # 选择数据中的最大日期作为基准,也可替换为F.current_date()获取当前日期 base_date = df_clean.select(F.max("dt")).first()[0] # 按类型分组,批量计算各窗口均值 result_df = df_clean.groupBy("ttype") \ .agg(*[ F.avg(F.when(F.datediff(base_date, F.col("dt")) <= day, F.col("amt"))).alias(f"avg_amt_{day}d") for day in window_days ]) result_df.show()
运行结果:
+-----+-----------+-----------+-----------+ |ttype|avg_amt_30d|avg_amt_60d|avg_amt_90d| +-----+-----------+-----------+-----------+ | C| 16.99| 14.99| 14.99| | D| 21.99| 21.99| 21.99| +-----+-----------+-----------+-----------+
步骤3:批量计算滚动时间窗口的均值(每条记录作为当前日期)
如果需求是每条记录对应的过去N天内同类型金额均值,可以用窗口函数结合列表推导式实现:
from pyspark.sql.window import Window # 定义窗口:按类型分区,按日期时间戳排序 window_spec = Window.partitionBy("ttype").orderBy(F.unix_timestamp("dt")) # 批量生成各滚动窗口均值列 rolling_result_df = df_clean.withColumn("dt_unix", F.unix_timestamp("dt")) \ .select("*", *[ F.avg("amt").over(window_spec.rangeBetween(-day*86400, 0)).alias(f"rolling_avg_{day}d") for day in window_days ]) \ .drop("dt_unix") rolling_result_df.show()
运行结果:
+-----+-----+----------+---------------+---------------+---------------+ |ttype| amt| dt|rolling_avg_30d|rolling_avg_60d|rolling_avg_90d| +-----+-----+----------+---------------+---------------+---------------+ | C|12.99|2024-01-01| 12.99| 12.99| 12.99| | C|16.99|2024-01-21| 14.99| 14.99| 14.99| | D|21.99|2024-02-15| 21.99| 21.99| 21.99| +-----+-----+----------+---------------+---------------+---------------+
优势说明
- 扩展性强:新增时间窗口只需修改
window_days列表,无需重复编写聚合逻辑 - 代码简洁:用列表推导式批量生成计算列,避免冗余代码
- 逻辑清晰:核心计算逻辑集中,可读性高
内容的提问来源于stack exchange,提问作者mithun_daa
相关产品推荐
相关产品推荐

