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

PySpark DataFrame滚动累计求和求助:输入转日期列累计输出失败

PySpark累计求和与日期列生成逻辑修正

现有输入DataFrame包含origin、destination以及「10+天」「10天」至「1天」的天数统计列,需求是从「10+天」列开始进行滚动累计求和,并将原天数列替换为当日及后续10天的日期列(格式为dd-MMM),生成目标输出DataFrame。但编写的PySpark代码未得到预期结果,需修正累计求和逻辑与日期生成逻辑。

输入DataFrame

origindestination10+天10天9天8天7天6天5天4天3天2天1天
CWCJMDCC660002113811263
CWCJPPSP2100021138823
PCWDMDCC500000000000
PCWDPPSP000000003039
DPMTJNPT000000000021
PMKMPPSP00000000000
PMKMMDCC20000000000

目标输出DataFrame

origindestination10+days8-Aug9-Aug10-Aug11-Aug12-Aug13-Aug14-Aug15-Aug16-Aug17-Aug18-Aug
CWCJMDCC666666666668698290101103166
CWCJPPSP212121212123243745535558
PCWDMDCC505050505050505050505050
PCWDPPSP0000000003342
DPMTJNPT0000000000021
PMKMPPSP000000000000
PMKMMDCC222222222222

错误代码

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()

当前错误输出

origindestination10+09-Aug10-Aug11-Aug12-Aug13-Aug14-Aug15-Aug16-Aug17-Aug18-Aug19-Aug
CWCJMDCC6633333333333
CWCJPPSP2100000000000
DPMTJNPT000000000000
PCWDMDCC504242424242424242424242
PCWDPPSP06363636363636363636363
PMKMMDCC200000000000
PMKMPPSP000000000000
Total139126126126126126126126126126126126

修正后的代码及逻辑说明

原代码存在三个核心问题:

  1. 窗口函数未按origin和destination分组,导致跨行求和
  2. 累计求和逻辑未实现从「10+天」开始的滚动累加
  3. 日期与天数列的映射关系错误

修正后的代码:

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.13 18:12:01