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

PySpark如何按小时聚合补全缺失小时并计算分组累计和

PySpark 按组补全全天24小时统计滚动累计值实现

问题说明

现有如下Spark DataFrame:

GroupIdEvent_timeEvent_nameEvent_value
xx2011-08-15 14:47:02.617023eventA1
xx2011-08-15 14:48:02.507053eventA2
xx2011-08-15 16:47:02.512016eventA100
yy2011-08-15 11:47:02.337019eventA2
yy2011-08-15 12:47:02.617041eventA1
yy2011-08-15 13:47:02.927040eventA3

需求为基于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处理)

预期输出结构如下(省略部分小时条目):

GroupIdDateHourCountagg_count
xx2011-08-1500NA0
xx2011-08-1501NA0
xx2011-08-1502NA0
xx2011-08-1513NA0
xx2011-08-151433
xx2011-08-1515NA3
xx2011-08-1516100103
xx2011-08-1517NA103
xx2011-08-1523NA103

原有实现代码如下,问题为无法补全单日全部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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.31 12:21:19