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

PySpark DataFrame透视转换及向前2行滚动求和实现方法

实现逻辑

你需要按「逆透视转年份为行→筛选目标name→透视转name为列→窗口计算2年滚动和」的步骤实现,具体代码如下:

完整PySpark代码

from pyspark.sql import SparkSession
from pyspark.sql import functions as F
from pyspark.sql.window import Window

# 初始化SparkSession,你也可以直接用你现有环境的SparkSession
spark = SparkSession.builder.appName("year_roll_sum").getOrCreate()

# ---------------------- 以下是构造测试原始数据,你可以替换为自己的DataFrame ----------------------
test_data = [
    ("abc", 100, 300),
    ("bbc", 200, 400),
    ("cbc", 300, 500),
    ("xyz", 700, 500),
    ("xzz", 200, 500)
]
original_df = spark.createDataFrame(test_data, schema=["name", "1990", "1991"])
# -----------------------------------------------------------------------------------------

# 1. 逆透视:把年份列转为行,生成year、value字段
unpivot_df = original_df.select(
    "name",
    # 这里stack的第一个参数是年份列的数量,有多少年就填多少,后面跟着 年份字符串, 对应年份列 成对拼接
    F.expr("stack(2, '1990', `1990`, '1991', `1991`) as (year, value)")
)

# 2. 筛选需要保留的name,不需要的话可以删掉这一步
filter_df = unpivot_df.filter(F.col("name").isin("abc", "bbc"))

# 3. 透视:把name字段取值转为列
pivot_df = filter_df.groupBy("year").pivot("name").agg(F.first("value"))

# 4. 定义2年滚动窗口:按年份升序排序,窗口范围是前1行到当前行,正好覆盖2年数据
roll_window = Window.orderBy("year").rowsBetween(-1, 0)

# 5. 计算各name的滚动和
result_df = pivot_df.select(
    "year",
    F.sum("abc").over(roll_window).alias("abc"),
    F.sum("bbc").over(roll_window).alias("bbc")
)

# 输出结果
result_df.show()

输出结果说明

运行上述代码后输出的结果如下,你问题描述中的100+200是累加示意,实际会返回真实的累加数值:

+----+---+---+
|year|abc|bbc|
+----+---+---+
|1990|100|200|
|1991|400|600|
+----+---+---+

如果需要保留所有name,删除筛选步骤,透视后批量对所有name列应用滚动求和逻辑即可。

内容的提问来源于stack exchange,提问作者ZZzzZZzz

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.03 20:57:03