PySpark中DataFrame转置需求:按ID分配专属列存储对应Code值
PySpark实现按ID分组分配列的需求
要实现将每个ID对应的code分配到连续的专属列中(比如R101占col1/col2,R201占col3/col4等),可以通过计算列位置+Pivot的方式来完成,下面是完整的解决方案:
步骤1:初始化环境与测试数据
首先创建测试DataFrame模拟你的输入数据:
from pyspark.sql import SparkSession from pyspark.sql import functions as F from pyspark.sql.window import Window # 初始化SparkSession spark = SparkSession.builder.appName("ID_Code_Column_Assign").getOrCreate() # 输入数据 input_data = [ ("R101", "GTR001"), ("R201", "RTY987"), ("R301", "KIT158"), ("R201", "PLI564"), ("R101", "MJU098"), ("R301", "OUY579") ] df = spark.createDataFrame(input_data, ["id", "code"]) df.show()
步骤2:计算每个ID的code数量与起始列位置
我们需要先确定每个ID需要占用多少列,以及这些列的起始位置:
# 计算每个ID的code总数 id_code_count = df.groupBy("id").agg(F.count("code").alias("code_count")) # 按ID排序,计算每个ID的起始列号(比如R101从col1开始,R201从col3开始) window_cumulative = Window.orderBy("id") id_col_mapping = id_code_count.withColumn( "start_col", F.sum("code_count").over(window_cumulative) - F.col("code_count") + 1 ) id_col_mapping.show()
步骤3:给每个code分配对应的目标列名
将起始列信息关联回原DataFrame,再给每个ID内的code分配具体的列名:
# 关联起始列信息到原数据 df_joined = df.join(id_col_mapping, on="id", how="left") # 给每个ID内的code分配序号(如果需要保持输入顺序,可改用monotonically_increasing_id()排序) window_code_order = Window.partitionBy("id").orderBy("code") df_with_col_info = df_joined.withColumn( "code_idx", F.row_number().over(window_code_order) ).withColumn( "col_num", F.col("start_col") + F.col("code_idx") - 1 ).withColumn( "col_name", F.concat(F.lit("col"), F.col("col_num").cast("string")) ) df_with_col_info.show()
步骤4:通过Pivot转换为目标结构
最后用pivot将行转列,得到你想要的输出:
# 执行Pivot操作,聚合得到每个列的code值 result_df = df_with_col_info.groupBy("id").pivot("col_name").agg(F.first("code")) result_df.show()
执行后你会得到如下结果:
+-----+-------+-------+-------+-------+-------+-------+ | id| col1| col2| col3| col4| col5| col6| +-----+-------+-------+-------+-------+-------+-------+ |R101 |GTR001 |MJU098 | null| null| null| null| |R201 | null| null|RTY987 |PLI564 | null| null| |R301 | null| null| null| null|KIT158 |OUY579 | +-----+-------+-------+-------+-------+-------+-------+
注意事项
- 如果需要保持
code的原始输入顺序(而非按code字典序排序),可以将window_code_order中的orderBy("code")替换为orderBy(F.monotonically_increasing_id()),因为monotonically_increasing_id()会保留数据的输入顺序。 - 该方案支持任意数量的code(每个ID的code数不固定),列名会自动按ID顺序连续分配。
内容的提问来源于stack exchange,提问作者user8510536
相关产品推荐
相关产品推荐

