如何在PySpark中处理多文件实现列的近3天累计求和
问题背景
有多个单月CSV数据文件,需为每个文件新增一列计算某列的近3天累计求和。现有代码仅能处理单个文件,无法解决跨月场景(如计算2020-04-01的累计值需要用上月最后3天数据),要求不合并所有CSV,实现每月一次的任务运行。
现有处理代码:
grouped_df = df.groupBy("province_state","date").sum('confirmed','deaths','recovered','active').orderBy("province_state","date") window = Window.partitionBy("province_state").orderBy("date").rowsBetween(-3, 0) window_df = grouped_df.select("*",sum('sum(active)').over(window).alias("3_days_window_active")) year_month_df = window_df.select("*",date_format("date", "yyyy-MM").alias("YYYY-MM")) year_month_df.write.option("header", True) \ .partitionBy("province_state","YYYY-MM")\ .mode("overwrite")\ .csv("state")
示例数据:
CSV-1(3月数据)
Date value cum_value
2020-03-25 7 7
2020-03-26 1 8
2020-03-27 1 9
2020-03-28 2 4
2020-03-29 4 7
2020-03-30 2 8
2020-03-31 6 12CSV-2(4月数据)
Date value cum_value
2020-04-01 7 15
2020-04-02 11 24
2020-04-03 13 31
2020-04-04 21 45
2020-04-05 40 73
2020-04-06 29 90
2020-04-07 60 129
可行解决方案
1. 预存上月末尾3天聚合数据
每月运行任务时,先读取上月的聚合后数据(或直接读取上月原始CSV重新聚合),提取每个province_state的最后3天目标列(如sum(active))数据,保存为临时DataFrame。
示例代码片段:
# 读取上月处理后的聚合数据 last_month_df = spark.read.csv("state/province_state=*/YYYY-MM=2020-03", header=True) # 按省份分组,取每个组的最后3天数据 from pyspark.sql.functions import row_number, col window_last = Window.partitionBy("province_state").orderBy(col("date").desc()) last_3_days_df = last_month_df.withColumn("rn", row_number().over(window_last)) \ .filter(col("rn") <= 3) \ .drop("rn", "3_days_window_active", "YYYY-MM") \ .orderBy("province_state", "date")
2. 合并上月数据与当月数据计算窗口
读取当月原始CSV,完成groupBy聚合后,将上月最后3天的数据追加到当月聚合数据的前面,再应用窗口函数计算累计值,最后筛选出当月数据写入结果。
示例代码片段:
# 读取并聚合当月数据 current_month_df = spark.read.csv("current_month_data.csv", header=True) grouped_current_month_df = current_month_df.groupBy("province_state","date") \ .sum('confirmed','deaths','recovered','active') \ .orderBy("province_state","date") # 合并上月最后3天数据与当月聚合数据 combined_df = last_3_days_df.union(grouped_current_month_df) # 计算3天窗口累计值 window = Window.partitionBy("province_state").orderBy("date").rowsBetween(-3, 0) window_df = combined_df.select("*", sum('sum(active)').over(window).alias("3_days_window_active")) # 筛选当月数据并写入 current_month = "2020-04" year_month_df = window_df.select("*",date_format("date", "yyyy-MM").alias("YYYY-MM")) \ .filter(col("YYYY-MM") == current_month) year_month_df.write.option("header", True) \ .partitionBy("province_state","YYYY-MM")\ .mode("overwrite")\ .csv("state")
3. 持久化上月末尾数据供下月使用
每次处理完当月数据后,提取每个province_state的最后3天聚合数据,单独存储到固定路径(如last_3_days目录),供下月任务直接读取,避免重复处理整月数据。
示例代码片段:
# 提取当月最后3天数据 current_last_3_days = year_month_df.withColumn("rn", row_number().over(window_last)) \ .filter(col("rn") <= 3) \ .drop("rn", "3_days_window_active", "YYYY-MM") # 写入持久化路径 current_last_3_days.write.option("header", True) \ .partitionBy("province_state")\ .mode("overwrite")\ .csv("last_3_days")
4. 边界场景处理
- 对于年度第一个月(如2020-01),无上月数据,直接用当月数据计算窗口即可,窗口函数会自动处理前几天的累计(如2020-01-01的累计值就是当天自身值)。
- 确保所有日期字段格式统一为
yyyy-MM-dd,避免排序错误导致窗口计算失效。
内容的提问来源于stack exchange,提问作者Surya Vamsi

