如何使用Python在Spark中无需聚合完成DataFrame转置操作
Spark Python 无聚合DataFrame转置实现方案
以下是针对你给出的输入输出需求的完整实现,因每个转置后的单元格映射唯一,全程无实际聚合计算:
实现思路
- 首先将原表除
COLUMN_NAME外的所有VALUE列拆分为长表结构,保留原VALUE列名作为转置后的行标识 - 基于
COLUMN_NAME的取值做列旋转(pivot),因每个行标识+列名组合仅对应唯一值,仅需用first取唯一值即可,无聚合逻辑 - 最终输出转置后的宽表结构
完整代码示例
from pyspark.sql import SparkSession from pyspark.sql.functions import first, expr # 初始化SparkSession spark = SparkSession.builder.appName("transpose_df_no_agg").getOrCreate() # ---------------------- 构造输入样例数据 ---------------------- data = [ ("col1", "val11", "val21", "val31", "valN1"), ("col2", "val12", "val22", "val32", "valN2"), ("col3", "val13", "val23", "val33", "valN3"), ("col4", "val14", "val24", "val34", "valN4"), ("col5", "val15", "val25", "val35", "valN5") ] columns = ["COLUMN_NAME", "VALUE1", "VALUE2", "VALUE3", "VALUEN"] df = spark.createDataFrame(data, schema=columns) df.show() # 输出输入表结构,和你给出的样例一致 # ---------------------- 核心转置逻辑 ---------------------- # 提取数值列(排除COLUMN_NAME的所有列) value_cols = [c for c in df.columns if c != "COLUMN_NAME"] # 构造stack表达式:n个列对应 "col1 val1 col2 val2 ..." 格式,n是value_cols的长度 stack_expr = f"stack({len(value_cols)}, {','.join([f'\'{c}\', {c}' for c in value_cols])}) as (original_row, val)" # 转成长表 long_df = df.select("COLUMN_NAME", expr(stack_expr)) # pivot转置,first仅取唯一值,无实际聚合 transposed_df = long_df.groupBy("original_row") \ .pivot("COLUMN_NAME") \ .agg(first("val")) \ .drop("original_row") # 可选择是否保留原来的VALUE列名标识 # 输出转置结果 transposed_df.show()
输出结果验证
运行后输出的转置表和你期望的结构完全一致:
+-----+-----+-----+-----+-----+ | col1| col2| col3| col4| col5| +-----+-----+-----+-----+-----+ |val11|val12|val13|val14|val15| |val21|val22|val23|val24|val25| |val31|val32|val33|val34|val35| |valN1|valN2|valN3|valN4|valN5| +-----+-----+-----+-----+-----+
内容的提问来源于stack exchange,提问作者Abhy
相关产品推荐
相关产品推荐

