在PySpark中实现按unique_id将固定行数转固定列数的行转列操作
PySpark实现多行转宽表(按unique_id聚合最多3行)
实现步骤:
- 添加分组内序号:用窗口函数
row_number()给每个unique_id下的行分配1-3的序号,可按需调整排序规则。 - 行转列(Pivot):以序号列为 pivot 列,将
column_A和column_B分别转成多列。 - 调整列名:把pivot生成的默认列名改成
column_A_1、column_B_2这类目标格式。
完整代码示例:
from pyspark.sql import SparkSession from pyspark.sql.window import Window from pyspark.sql.functions import row_number, first, col # 初始化SparkSession spark = SparkSession.builder.appName("row_to_wide").getOrCreate() # 模拟输入数据 data = [ (123, 12345, "ABCDEFG"), (123, 23456, "BCDEFGH"), (123, 34567, "CDEFGHI"), (234, 12345, "ABCDEFG") ] df = spark.createDataFrame(data, ["unique_id", "column_A", "column_B"]) # 步骤1:添加分组内序号 window_spec = Window.partitionBy("unique_id").orderBy("column_A") # 可修改排序字段适配业务 df_with_rank = df.withColumn("rank", row_number().over(window_spec)) # 步骤2:Pivot行转列 pivoted_df = df_with_rank.groupBy("unique_id") \ .pivot("rank", [1, 2, 3]) \ .agg(first("column_A").alias("A"), first("column_B").alias("B")) # 步骤3:调整列名到目标格式 final_df = pivoted_df.select( "unique_id", col("1_A").alias("column_A_1"), col("2_A").alias("column_A_2"), col("3_A").alias("column_A_3"), col("1_B").alias("column_B_1"), col("2_B").alias("column_B_2"), col("3_B").alias("column_B_3") ) # 输出结果 final_df.show()
关键说明:
- 窗口函数的
orderBy可根据业务需求替换(比如按时间字段排序),确保序号分配符合预期。 pivot时指定[1,2,3]作为固定序号值,保证即使某个unique_id不足3行,也会生成对应的空列。- 聚合函数用
first()是因为每个序号在分组内唯一,取第一个值即可,也可以用max()等等价函数。
内容的提问来源于stack exchange,提问作者Robert Kossendey
相关产品推荐
相关产品推荐

