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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.06 18:54:58