PySpark如何按小时聚合补全缺失小时并计算分组累计和
PySpark 按组补全全天24小时统计滚动累计值实现
问题说明
现有如下Spark DataFrame:
| GroupId | Event_time | Event_name | Event_value |
|---|---|---|---|
| xx | 2011-08-15 14:47:02.617023 | eventA | 1 |
| xx | 2011-08-15 14:48:02.507053 | eventA | 2 |
| xx | 2011-08-15 16:47:02.512016 | eventA | 100 |
| yy | 2011-08-15 11:47:02.337019 | eventA | 2 |
| yy | 2011-08-15 12:47:02.617041 | eventA | 1 |
| yy | 2011-08-15 13:47:02.927040 | eventA | 3 |
需求为基于GroupId分组,统计单日维度下每小时的eventA对应Event_value滚动累计值:
- 对
GroupId为xx、时间为2011-08-15 14:00的条目,需计算该GroupId下14:00到15:00区间内eventA的Event_value之和,本场景计算结果为1+2=3 - 需展示单日00点到23点的全部小时条目,若对应小时区间无eventA事件,则该小时的Count值记为NA(后续计算按0处理)
预期输出结构如下(省略部分小时条目):
| GroupId | Date | Hour | Count | agg_count |
|---|---|---|---|---|
| xx | 2011-08-15 | 00 | NA | 0 |
| xx | 2011-08-15 | 01 | NA | 0 |
| xx | 2011-08-15 | 02 | NA | 0 |
| xx | 2011-08-15 | 13 | NA | 0 |
| xx | 2011-08-15 | 14 | 3 | 3 |
| xx | 2011-08-15 | 15 | NA | 3 |
| xx | 2011-08-15 | 16 | 100 | 103 |
| xx | 2011-08-15 | 17 | NA | 103 |
| xx | 2011-08-15 | 23 | NA | 103 |
原有实现代码如下,问题为无法补全单日全部24小时的条目,无法输出符合要求的结果:
from pyspark.sql.functions import col, count, hour, sum df2 = (df .withColumn("Event_time", col("Event_time").cast("timestamp")) .withColumn("Date", col("Event_time").cast("date")) .withColumn("Hour", hour(col("Event_time")))) df3 = df2.groupBy("GroupId", "Date", "Hour").count() df3.withColumn( "agg_count", sum("Count").over(Window.partitionBy("GroupId", "Date").orderBy("Hour")))
调整方案
原有代码聚合后仅保留了存在事件的小时记录,需要先构造所有GroupId+Date+0~23小时的全量组合,再和聚合结果关联补全数据,最后计算滚动累计值。完整实现代码如下:
from pyspark.sql.functions import col, hour, sum as spark_sum, lit, sequence, explode, to_date, when from pyspark.sql.window import Window # 1. 基础字段转换,过滤非eventA事件 base_df = (df .filter(col("Event_name") == "eventA") .withColumn("Event_time", col("Event_time").cast("timestamp")) .withColumn("Date", to_date(col("Event_time"))) .withColumn("Hour", hour(col("Event_time"))) ) # 2. 按组、日期、小时聚合Event_value总和 hourly_agg_df = (base_df .groupBy("GroupId", "Date", "Hour") .agg(spark_sum("Event_value").alias("Count")) ) # 3. 构造全量小时维度:获取所有(GroupId, Date)组合,生成0-23点的完整小时序列 group_date_df = base_df.select("GroupId", "Date").distinct() full_hour_df = (group_date_df .withColumn("Hour", explode(sequence(lit(0), lit(23)))) ) # 4. 关联聚合结果,补全缺失小时的Count值为null(对应NA),null转0用于计算累计值 full_df = (full_hour_df .join(hourly_agg_df, on=["GroupId", "Date", "Hour"], how="left") .withColumn("Count_calc", when(col("Count").isNull(), 0).otherwise(col("Count"))) ) # 5. 定义窗口计算滚动累计值,最后将Count为null的替换为NA window_spec = Window.partitionBy("GroupId", "Date").orderBy("Hour").rowsBetween(Window.unboundedPreceding, Window.currentRow) result_df = (full_df .withColumn("agg_count", spark_sum("Count_calc").over(window_spec)) .withColumn("Count", when(col("Count").isNull(), "NA").otherwise(col("Count").cast("string"))) .select("GroupId", "Date", "Hour", "Count", "agg_count") .orderBy("GroupId", "Date", "Hour") ) result_df.show()
关键逻辑说明
- 提前过滤只保留
eventA事件,避免其他事件干扰统计结果 - 通过
sequence(0,23)生成0到23的小时序列,和所有GroupId+Date组合做笛卡尔积,得到全量小时条目 - 左关联聚合结果后,无事件的小时Count为null,计算累计值时按0处理,最终展示时替换为NA
- 窗口函数设置为按组、日期分区,按小时排序,从分区第一行累计到当前行,得到符合要求的滚动累计值
内容的提问来源于stack exchange,提问作者user18707304
相关产品推荐
相关产品推荐

