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
相关产品推荐
相关产品推荐

