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

PySpark按row_number透视表结果不符,求代码问题排查

问题分析与解决方法

你的代码问题出在pivot的字段选错了:你用pivot("exploded_score"),Spark会把exploded_score的每个唯一数值直接作为列名,这就是为什么结果列名是数值的原因。而你需要的是按row_number分组后,把每组内的exploded_score按顺序拆成score_1、score_2这类序号化的列,核心是要先给每组内的exploded_score分配一个位置序号,再透视这个序号字段。

正确实现步骤

  1. 给每组内的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))
    
  2. 透视位置序号字段
    现在用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"))
    
  3. 可选:重命名列名
    如果想把列名从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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.18 01:15:06