PySpark长表转宽表:拆分amount列并自定义新列名方法问询
PySpark实现长表转宽并拆分amount列
步骤1:给分组内的行添加序号
首先通过窗口函数给每个(id, term)组内的amount行分配唯一序号,为后续拆分列做准备:
from pyspark.sql import Window from pyspark.sql.functions import row_number # 定义窗口规则:按id、term分组,可按需调整排序逻辑 window_spec = Window.partitionBy("id", "term").orderBy("amount") # 添加序号列rank df_with_rank = df.withColumn("rank", row_number().over(window_spec)) df_with_rank.show()
输出结果:
+---+----+------+----+ | id|term|amount|rank| +---+----+------+----+ | 1| 6| 200| 1| | 1| 6| 300| 2| | 1| 6| 400| 3| | 1| 7| 1000| 1| | 1| 7| 5000| 2| +---+----+------+----+
步骤2:转宽并自定义列名
通过pivot转置序号列,再将默认列名重命名为自定义格式:
# 分组后转置rank列,聚合取唯一amount值 wide_df = df_with_rank.groupBy("id", "term") \ .pivot("rank") \ .agg({"amount": "first"}) \ .withColumnRenamed("1", "amount_1") \ .withColumnRenamed("2", "amount_2") \ .withColumnRenamed("3", "amount_3") wide_df.show()
最终结果:
+---+----+--------+--------+--------+ | id|term|amount_1|amount_2|amount_3| +---+----+--------+--------+--------+ | 1| 6| 200| 300| 400| | 1| 7| 1000| 5000| null| +---+----+--------+--------+--------+
动态适配任意分组行数
如果不同分组的amount行数不固定,可通过动态方式生成列名,避免手动逐个重命名:
# 提取所有出现过的rank值,生成列名映射 rank_values = [row.rank for row in df_with_rank.select("rank").distinct().orderBy("rank").collect()] col_mapping = {str(r): f"amount_{r}" for r in rank_values} # 执行转置 wide_df = df_with_rank.groupBy("id", "term").pivot("rank").agg({"amount": "first"}) # 批量重命名列 for old_col, new_col in col_mapping.items(): wide_df = wide_df.withColumnRenamed(old_col, new_col) wide_df.show()
内容的提问来源于stack exchange,提问作者Chuck
相关产品推荐
相关产品推荐

