PySpark中自定义列名实现表透视(Pivot)需求求助
PySpark实现透视表并自定义列名
要实现你需要的转换效果,核心是先给分组内的记录添加序号,再通过透视将序号映射为自定义列名,具体步骤如下:
1. 准备测试数据
先模拟你的原始数据表:
from pyspark.sql import SparkSession from pyspark.sql.window import Window from pyspark.sql.functions import row_number, first, col spark = SparkSession.builder.appName("PivotCustomCol").getOrCreate() data = [ (100, "2022Q2", "202204", "abc"), (200, "2022Q2", "202204", "abc"), (300, "2022Q2", "202204", "abc") ] df = spark.createDataFrame(data, ["balance", "ex_cy", "rp_prd", "scenario"]) df.show()
2. 给分组内记录添加序号
用窗口函数,以ex_cy、rp_prd、scenario作为分组依据,给每组内的记录生成递增序号,后续用这个序号做透视:
window_spec = Window.partitionBy("ex_cy", "rp_prd", "scenario").orderBy("balance") df_with_rn = df.withColumn("rn", row_number().over(window_spec)) df_with_rn.show()
这里orderBy("balance")是为了保证序号和balance的数值顺序对应,如果你不需要按balance排序,可以替换为其他业务字段,或者去掉orderBy(但Spark会随机排序,建议保留排序逻辑保证结果稳定)。
3. 透视并重命名列
先按分组字段聚合,透视序号列,再将生成的数字列重命名为C1、C2、C3:
pivoted_df = df_with_rn.groupBy("ex_cy", "rp_prd", "scenario") \ .pivot("rn") \ .agg(first("balance")) \ .withColumnsRenamed({ "1": "C1", "2": "C2", "3": "C3" }) pivoted_df.show()
这里用first("balance")是因为每个序号在分组内唯一,使用sum、max等聚合函数结果一致,可根据习惯选择。
最终输出结果
执行后会得到你需要的表结构:
+------+------+--------+---+---+---+ | ex_cy|rp_prd|scenario| C1| C2| C3| +------+------+--------+---+---+---+ |2022Q2|202204| abc|100|200|300| +------+------+--------+---+---+---+
额外说明
如果分组内的记录数量不固定(比如有的组有5条,有的组有3条),可以通过遍历透视后的列名动态生成重命名规则,示例代码如下:
# 动态重命名所有数字列 new_col_names = [col(c).alias(f"C{c}") if c.isdigit() else col(c) for c in pivoted_df.columns] pivoted_df = pivoted_df.select(new_col_names)
内容的提问来源于stack exchange,提问作者Mohammad Sunny
相关产品推荐
相关产品推荐

