如何在PySpark中按天对1小时窗口进行累计聚合
解决Spark DataFrame按group_id统计每日小时累计事件数的问题
我来帮你搞定这个累计统计的需求!你想要的是按group_id统计从当天00:00到每个小时结束的累计事件总数,而不是单小时的独立计数,咱们一步步来实现:
步骤1:提取时间维度字段
首先,我们需要从event_time里拆分出日期和小时,方便后续按天和小时分组:
from pyspark.sql import functions as F from pyspark.sql.window import Window # 提取日期、小时字段,同时保留原始时间字段 df_with_time = df.withColumn("event_date", F.to_date(F.col("event_time"))) \ .withColumn("hour_of_day", F.hour(F.col("event_time")))
步骤2:补全所有小时(处理无事件的小时)
你的预期输出里包含了事件数为0的小时(比如YYYY在01点窗口的计数),所以我们需要生成每个group_id每天的全部24小时,避免缺失:
# 获取所有唯一的group_id和日期组合 unique_groups_dates = df_with_time.select("group_id", "event_date").distinct() # 生成0到23的小时序列 hours = spark.range(0, 24).withColumnRenamed("id", "hour_of_day") # 通过笛卡尔积得到每个group_id+日期对应的所有小时 all_hours = unique_groups_dates.crossJoin(hours)
步骤3:统计单小时事件数并补全0
把原始数据的单小时计数和全量小时表左连接,没有事件的小时填充0:
# 统计每个group_id+日期+小时的事件数 hourly_counts = df_with_time.groupBy("group_id", "event_date", "hour_of_day") \ .agg(F.count("*").alias("hourly_count")) # 左连接补全缺失小时,无事件的小时计数设为0 full_hourly_counts = all_hours.join(hourly_counts, ["group_id", "event_date", "hour_of_day"], how="left") \ .fillna(0, subset=["hourly_count"])
步骤4:计算累计事件数
用窗口函数实现从当天0点到当前小时的累计求和:
# 定义窗口规则:按group_id和日期分区,按小时升序排序,累加从第一行到当前行的所有计数 cumulative_window = Window.partitionBy("group_id", "event_date") \ .orderBy("hour_of_day") \ .rowsBetween(Window.unboundedPreceding, Window.currentRow) # 计算累计值 cumulative_df = full_hourly_counts.withColumn("values", F.sum("hourly_count").over(cumulative_window))
步骤5:构造预期的窗口格式
最后把时间转换成你需要的[开始时间,结束时间]数组格式:
# 构造窗口的开始(当天00:00)和结束(当前小时+1点)时间 result_df = cumulative_df.withColumn("window_start", F.to_timestamp(F.concat(F.col("event_date"), F.lit(" 00:00:00")))) \ .withColumn("window_end", F.to_timestamp(F.concat(F.col("event_date"), F.lit(" "), F.col("hour_of_day")+1, F.lit(":00:00")))) \ .withColumn("model_window", F.array("window_start", "window_end")) \ .select("group_id", "model_window", "values") # 按group_id和窗口结束时间排序展示 result_df.orderBy("group_id", "window_end").show(truncate=False)
为什么你的原始代码不行?
你之前用的window('event_time', '1 hour')是滚动窗口,每个窗口只统计当前1小时内的事件数,而不是从当天起始时间的累计值。我们上面的方案先拆分时间维度、补全所有小时,再用累计窗口函数,完美匹配你的需求。
内容的提问来源于stack exchange,提问作者Stergios
相关产品推荐
相关产品推荐

