PySpark按row_number透视表结果不符,求代码问题排查
你的代码问题出在pivot的字段选错了:你用pivot("exploded_score"),Spark会把exploded_score的每个唯一数值直接作为列名,这就是为什么结果列名是数值的原因。而你需要的是按row_number分组后,把每组内的exploded_score按顺序拆成score_1、score_2这类序号化的列,核心是要先给每组内的exploded_score分配一个位置序号,再透视这个序号字段。
正确实现步骤
给每组内的
exploded_score添加位置序号
用窗口函数,按你需要保留的维度(比如index、Class、PBUP_AC、item_code、Running_Index、row_number)分组,给每个exploded_score分配一个顺序编号(比如pos)。如果需要指定排序规则,在orderBy里加对应的字段,没有的话可以用monotonically_increasing_id()保证顺序:from pyspark.sql import Window from pyspark.sql.functions import row_number, first, monotonically_increasing_id # 定义窗口:按需要保留的字段分区,按指定规则排序(这里用monotonically_increasing_id()保证原有顺序) window_spec = Window.partitionBy("index", "Class", "PBUP_AC", "item_code", "Running_Index", "row_number") \ .orderBy(monotonically_increasing_id()) # 添加位置序号列pos df_with_pos = exploded_df_1.withColumn("pos", row_number().over(window_spec))透视位置序号字段
现在用pos字段做pivot,而不是exploded_score,这样列名会是1、2、3...,再聚合取对应的exploded_score值:# 按维度字段分组,透视pos列 pivoted_df = df_with_pos.groupBy("index", "Class", "PBUP_AC", "item_code", "Running_Index", "row_number") \ .pivot("pos") \ .agg(first("exploded_score"))可选:重命名列名
如果想把列名从1、2改成score_1、score_2这类更直观的名称:# 遍历列名,对数字列重命名 new_columns = [f"score_{col}" if col.isdigit() else col for col in pivoted_df.columns] pivoted_df = pivoted_df.toDF(*new_columns)
为什么原来的代码不对?
你原来的pivot("exploded_score")是把exploded_score的每个唯一值作为列的标识,比如如果某个row_number下有exploded_score为5、8、10,就会生成列名5、8、10的列,这和你需要的“按顺序生成独立列”的需求完全不符。必须先给每组内的exploded_score分配顺序序号,再透视这个序号才能得到预期结果。
内容的提问来源于stack exchange,提问作者Malek BEN HMIDA

