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

如何在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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.15 08:35:59