PySpark中如何基于15分钟滚动窗口计算分组内子组平均均值?
解决Spark 15分钟滚动窗口计算城市均值的平均问题
核心思路
要实现过去15分钟内各城市销售均值的平均值(每个城市权重平等,不受记录数量影响),需要分两步处理:
- 先计算每个城市在过去15分钟内的滚动销售均值
- 针对每一行的时间,聚合同区域内所有城市的最新滚动均值,再取其平均
具体实现步骤
1. 预处理时间列
先将字符串格式的时间转换为Spark可识别的Timestamp类型:
from pyspark.sql import functions as F from pyspark.sql.window import Window # 转换时间列格式 df = df.withColumn("Datetime", F.to_timestamp("Datetime"))
2. 计算城市级15分钟滚动均值
创建按Country+Region+City分区的时间窗口,计算每个城市在当前时间往前15分钟内的销售均值:
# 定义城市级滚动窗口:15分钟转换为秒(900秒) city_window = Window.partitionBy("Country", "Region", "City") \ .orderBy(F.unix_timestamp("Datetime")) \ .rangeBetween(-15*60, 0) # 新增城市滚动均值列 df_with_city_avg = df.withColumn("city_avg", F.avg("Sale").over(city_window))
3. 获取每个城市的最新滚动均值
同一城市可能有多条记录落在15分钟窗口内,我们需要保留每个城市截至当前时间的最新均值:
# 定义窗口:按城市分区、时间排序,取最新的均值 latest_city_window = Window.partitionBy("Country", "Region", "City") \ .orderBy(F.unix_timestamp("Datetime")) df_with_latest_avg = df_with_city_avg.withColumn("latest_city_avg", F.last("city_avg").over(latest_city_window))
4. 计算区域级的均值平均
创建按Country+Region分区的时间窗口,收集窗口内所有城市的最新均值,再计算其平均:
# 定义区域级滚动窗口 region_window = Window.partitionBy("Country", "Region") \ .orderBy(F.unix_timestamp("Datetime")) \ .rangeBetween(-15*60, 0) # 收集窗口内的城市-最新均值集合,再计算平均 df_result = df_with_latest_avg.withColumn("city_avg_collection", F.collect_set(F.struct("City", "latest_city_avg")).over(region_window)) \ .withColumn("region_sale_average", F.expr(""" aggregate( city_avg_collection, 0.0, (acc, item) -> acc + item.latest_city_avg, acc -> acc / size(city_avg_collection) ) """))
结果验证
最终df_result中的region_sale_average列会得到预期值:100, 550, 547.5, 546,完全符合需求。
内容的提问来源于stack exchange,提问作者Jokin
相关产品推荐
相关产品推荐

