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

如何在PySpark中处理多文件实现列的近3天累计求和

跨月度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 12

CSV-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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.14 15:20:22