PySpark DataFrame滚动累计求和求助:输入转日期列累计输出失败
PySpark累计求和与日期列生成逻辑修正
现有输入DataFrame包含origin、destination以及「10+天」「10天」至「1天」的天数统计列,需求是从「10+天」列开始进行滚动累计求和,并将原天数列替换为当日及后续10天的日期列(格式为dd-MMM),生成目标输出DataFrame。但编写的PySpark代码未得到预期结果,需修正累计求和逻辑与日期生成逻辑。
输入DataFrame
| origin | destination | 10+天 | 10天 | 9天 | 8天 | 7天 | 6天 | 5天 | 4天 | 3天 | 2天 | 1天 |
|---|---|---|---|---|---|---|---|---|---|---|---|---|
| CWCJ | MDCC | 66 | 0 | 0 | 0 | 2 | 1 | 13 | 8 | 11 | 2 | 63 |
| CWCJ | PPSP | 21 | 0 | 0 | 0 | 2 | 1 | 13 | 8 | 8 | 2 | 3 |
| PCWD | MDCC | 50 | 0 | 0 | 0 | 0 | 0 | 0 | 0 | 0 | 0 | 0 |
| PCWD | PPSP | 0 | 0 | 0 | 0 | 0 | 0 | 0 | 0 | 3 | 0 | 39 |
| DPMT | JNPT | 0 | 0 | 0 | 0 | 0 | 0 | 0 | 0 | 0 | 0 | 21 |
| PMKM | PPSP | 0 | 0 | 0 | 0 | 0 | 0 | 0 | 0 | 0 | 0 | 0 |
| PMKM | MDCC | 2 | 0 | 0 | 0 | 0 | 0 | 0 | 0 | 0 | 0 | 0 |
目标输出DataFrame
| origin | destination | 10+days | 8-Aug | 9-Aug | 10-Aug | 11-Aug | 12-Aug | 13-Aug | 14-Aug | 15-Aug | 16-Aug | 17-Aug | 18-Aug |
|---|---|---|---|---|---|---|---|---|---|---|---|---|---|
| CWCJ | MDCC | 66 | 66 | 66 | 66 | 66 | 68 | 69 | 82 | 90 | 101 | 103 | 166 |
| CWCJ | PPSP | 21 | 21 | 21 | 21 | 21 | 23 | 24 | 37 | 45 | 53 | 55 | 58 |
| PCWD | MDCC | 50 | 50 | 50 | 50 | 50 | 50 | 50 | 50 | 50 | 50 | 50 | 50 |
| PCWD | PPSP | 0 | 0 | 0 | 0 | 0 | 0 | 0 | 0 | 0 | 3 | 3 | 42 |
| DPMT | JNPT | 0 | 0 | 0 | 0 | 0 | 0 | 0 | 0 | 0 | 0 | 0 | 21 |
| PMKM | PPSP | 0 | 0 | 0 | 0 | 0 | 0 | 0 | 0 | 0 | 0 | 0 | 0 |
| PMKM | MDCC | 2 | 2 | 2 | 2 | 2 | 2 | 2 | 2 | 2 | 2 | 2 | 2 |
错误代码
from pyspark.sql.window import Window from pyspark.sql import functions as F from datetime import datetime, timedelta todays_date = datetime.today().date() future_dates = [todays_date + timedelta(days=i) for i in range(1, 12)] columns_to_sum = ["9", "8", "7", "6", "5", "4", "3", "2", "1"] window_spec = Window.orderBy("origin", "destination") for i, date in enumerate(future_dates): col_name = date.strftime("%d-%b") for col in columns_to_sum: summary_export_dwell_df_temp = summary_export_dwell_df_temp.withColumn(col_name, F.when(F.col(col) > 0, F.sum(col).over(window_spec)).otherwise(0).cast("int")) window_spec = Window.orderBy("origin", "destination") selected_columns = ["origin", "destination", "10+"] + [date.strftime("%d-%b") for date in future_dates] selected_df = summary_export_dwell_df_temp.select(*selected_columns) selected_df.show()
当前错误输出
| origin | destination | 10+ | 09-Aug | 10-Aug | 11-Aug | 12-Aug | 13-Aug | 14-Aug | 15-Aug | 16-Aug | 17-Aug | 18-Aug | 19-Aug |
|---|---|---|---|---|---|---|---|---|---|---|---|---|---|
| CWCJ | MDCC | 66 | 3 | 3 | 3 | 3 | 3 | 3 | 3 | 3 | 3 | 3 | 3 |
| CWCJ | PPSP | 21 | 0 | 0 | 0 | 0 | 0 | 0 | 0 | 0 | 0 | 0 | 0 |
| DPMT | JNPT | 0 | 0 | 0 | 0 | 0 | 0 | 0 | 0 | 0 | 0 | 0 | 0 |
| PCWD | MDCC | 50 | 42 | 42 | 42 | 42 | 42 | 42 | 42 | 42 | 42 | 42 | 42 |
| PCWD | PPSP | 0 | 63 | 63 | 63 | 63 | 63 | 63 | 63 | 63 | 63 | 63 | 63 |
| PMKM | MDCC | 2 | 0 | 0 | 0 | 0 | 0 | 0 | 0 | 0 | 0 | 0 | 0 |
| PMKM | PPSP | 0 | 0 | 0 | 0 | 0 | 0 | 0 | 0 | 0 | 0 | 0 | 0 |
| Total | 139 | 126 | 126 | 126 | 126 | 126 | 126 | 126 | 126 | 126 | 126 | 126 |
修正后的代码及逻辑说明
原代码存在三个核心问题:
- 窗口函数未按
origin和destination分组,导致跨行求和 - 累计求和逻辑未实现从「10+天」开始的滚动累加
- 日期与天数列的映射关系错误
修正后的代码:
from pyspark.sql.window import Window from pyspark.sql import functions as F from datetime import datetime, timedelta # 获取当日及后续10天的日期(共11个日期,对应10+天 + 10天到1天的累计) todays_date = datetime.today().date() date_list = [todays_date + timedelta(days=i) for i in range(0, 11)] date_cols = [d.strftime("%d-%b") for d in date_list] # 定义原始天数列的顺序:从10天到1天,对应后续10个日期列 day_cols = ["10天", "9天", "8天", "7天", "6天", "5天", "4天", "3天", "2天", "1天"] # 按origin和destination分组,确保每行独立计算累计 window_spec = Window.partitionBy("origin", "destination").orderBy(F.lit(1)) # 先将10+天列重命名为目标列名 df = summary_export_dwell_df.withColumnRenamed("10+天", "10+days") # 计算累计求和的基础列:从10+天开始,依次累加后续天数 cumulative_cols = ["10+days"] for day_col in day_cols: prev_cum = cumulative_cols[-1] new_cum = F.col(prev_cum) + F.col(day_col) cumulative_cols.append(new_cum.alias(f"cum_{day_col}")) # 生成累计求和的中间列 df_cumulative = df.select("origin", "destination", *cumulative_cols) # 将累计列映射到对应的日期列 for i in range(len(date_cols)): df_cumulative = df_cumulative.withColumnRenamed( cumulative_cols[i], date_cols[i] if i > 0 else "10+days" ) # 调整列顺序,与目标输出一致 final_cols = ["origin", "destination", "10+days"] + date_cols[1:] final_df = df_cumulative.select(final_cols) final_df.show()
关键修正点:
- 窗口函数:使用
partitionBy("origin", "destination")确保每个origin-destination组内独立计算,避免跨行干扰 - 累计逻辑:从「10+天」开始,依次累加「10天」「9天」...「1天」的值,生成滚动累计序列
- 日期映射:将生成的累计列与当日及后续10天的日期一一对应,确保列名和数值匹配目标输出
内容的提问来源于stack exchange,提问作者vish
相关产品推荐
相关产品推荐

