PySpark添加分组计数列及7天滚动求和实现问询
解决方案
步骤1:添加单条记录的递增计数(Count列)
你需要用row_number()窗口函数,按社区(Community_Area)+ 天级日期分区,再按原始时间戳排序,实现同社区同一天内的计数递增:
from pyspark.sql import functions as F from pyspark.sql.window import Window # 提取Date字段的天级日期(用于按天分组计数) df_with_date_part = df4b.withColumn("date_only", F.to_date(F.col("Date"), "MM/dd/yyyy")) # 定义窗口:按社区+天级日期分区,按原始时间戳排序 count_window = Window.partitionBy("Community_Area", "date_only").orderBy("Date") # 添加Count列 df_with_count = df_with_date_part.withColumn("Count", F.row_number().over(count_window)) # 查看中间结果(可选) df_with_count.select("Community_Area", "Date", "Count").show()
步骤2:计算7天滚动求和
要实现7天滚动求和,需先聚合每日总次数,再用滑动窗口计算过去7天的总和:
2.1 聚合每日总次数
先统计每个社区每天的总记录数:
# 按社区和天级日期分组,统计当日总记录数 daily_counts = df_with_count.groupBy("Community_Area", "date_only")\ .agg(F.count("*").alias("daily_total"))\ .orderBy("Community_Area", "date_only")
2.2 计算7天滚动求和
定义滑动窗口:按社区分区,按日期排序,窗口范围覆盖当前日期及往前6天(共7天):
# 定义7天滚动窗口:将日期转为时间戳后,用秒数计算窗口范围(-6*86400表示往前推6天) rolling_window = Window.partitionBy("Community_Area")\ .orderBy(F.col("date_only").cast("timestamp"))\ .rangeBetween(-6*86400, 0) # 添加滚动求和列,整理最终输出字段 final_df = daily_counts.withColumn("Rolling_7-day_sum", F.sum("daily_total").over(rolling_window))\ .select("Community_Area", F.col("date_only").alias("Date"), "Rolling_7-day_sum") # 查看最终结果 final_df.show()
关键说明
row_number()保证同社区同一天内,记录按时间顺序从1开始递增计数,完全匹配你要的效果。- 使用
rangeBetween而非rowsBetween是因为要按日期范围计算,而非按行数,避免日期不连续时的错误。 - 如果你的
Date字段已经是Spark的日期/时间戳类型,可以跳过to_date的转换步骤。
内容的提问来源于stack exchange,提问作者slycooper
相关产品推荐
相关产品推荐

