You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

PySpark中如何基于15分钟滚动窗口计算分组内子组平均均值?

解决Spark 15分钟滚动窗口计算城市均值的平均问题

核心思路

要实现过去15分钟内各城市销售均值的平均值(每个城市权重平等,不受记录数量影响),需要分两步处理:

  1. 先计算每个城市在过去15分钟内的滚动销售均值
  2. 针对每一行的时间,聚合同区域内所有城市的最新滚动均值,再取其平均

具体实现步骤

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.16 09:33:16