PySpark窗口分区实现年度销售列转换去重求助
解决PySpark窗口分区后重复数据问题
你的代码使用窗口函数时,分区包含了calyear,导致每个(ZZA_TEXT4, calyear)组合的每一行都会重复显示该组的汇总值,最终同一个ZZA_TEXT4对应多行记录。要得到每个ZZA_TEXT4一行、年度销售额为单独列的结果,需要先聚合再透视,而非直接在原数据上用窗口函数。
解决方案步骤
- 按
ZZA_TEXT4和calyear分组,计算每组的年度销售额总和 - 对
calyear字段进行透视,将年度值转为列名,对应值为该年度的销售额 - (可选)将透视后的空值填充为0,并重命名列以匹配期望格式
代码示例
from pyspark.sql import functions as F # 1. 聚合年度销售额 aggregated_df = df.groupBy("ZZA_TEXT4", "calyear").agg( F.sum("net_sales_op_sum").alias("annual_sales") ) # 2. 透视转换为年度列 pivoted_df = aggregated_df.groupBy("ZZA_TEXT4").pivot("calyear").agg( F.first("annual_sales") ) # 3. 调整列名并填充空值为0 final_df = pivoted_df.select( "ZZA_TEXT4", F.col("2018").alias("SALE_2018_Sales_OP").fillna(0), # 若有其他年份,按此格式添加即可 # F.col("2019").alias("SALE_2019_Sales_OP").fillna(0), ) final_df.show()
这样处理后,每个ZZA_TEXT4只会保留一行记录,各年度销售额以单独列展示,完全匹配你期望的输出效果。
内容的提问来源于stack exchange,提问作者Arun Mohan
相关产品推荐
相关产品推荐

