如何在PySpark DataFrame中将多列层级数据转置为多行结构
PySpark宽表转长表(行列转置)实现方案
你需要的操作是将多列level字段从宽表转为行维度的长表,直接使用PySpark内置的stack函数即可实现,不需要自定义UDF。
完整实现代码
1. 核心转换逻辑(适配现有数据集)
from pyspark.sql.functions import expr # 假设你的原始DataFrame名为source_df result_df = source_df.select( "UNQ_ID", # stack参数说明:4=要转置的level列总数,后面每两组对应原列名、原列值,最后指定输出列别名 expr("stack(4, 'level_1', level_1, 'level_2', level_2, 'level_3', level_3, 'level_4', level_4) as (level, name)") ) # 可选:按UNQ_ID和level排序,得到和示例完全一致的输出顺序 result_df = result_df.orderBy("UNQ_ID", "level") # 查看结果 result_df.show()
2. 动态适配多level列场景(无需手动写全列名)
如果后续level字段数量变化,可自动识别所有level前缀的列生成转换逻辑:
from pyspark.sql.functions import expr # 自动识别所有level开头的待转置列 level_cols = [col for col in source_df.columns if col.startswith("level_")] # 动态生成stack表达式 stack_expr = f"stack({len(level_cols)}, {', '.join([f'\'{c}\', {c}' for c in level_cols])}) as (level, name)" # 执行转换 result_df = source_df.select("UNQ_ID", expr(stack_expr)).orderBy("UNQ_ID", "level")
3. 含测试数据的完整可运行示例
from pyspark.sql import SparkSession from pyspark.sql.functions import expr # 初始化SparkSession spark = SparkSession.builder.appName("transpose_demo").getOrCreate() # 构造和示例一致的测试数据 test_data = [ ("D Group", "Investments", "ORB", "ECM", 1), ("E Group", "Investment", "Origination", "Execution", 2) ] schema = ["level_1", "level_2", "level_3", "level_4", "UNQ_ID"] source_df = spark.createDataFrame(test_data, schema=schema) # 执行转置 level_cols = [col for col in source_df.columns if col.startswith("level_")] stack_expr = f"stack({len(level_cols)}, {', '.join([f'\'{c}\', {c}' for c in level_cols])}) as (level, name)" result_df = source_df.select("UNQ_ID", expr(stack_expr)).orderBy("UNQ_ID", "level") # 输出结果 result_df.show()
运行后输出和要求的目标结构完全一致。
内容的提问来源于stack exchange,提问作者santhosh
相关产品推荐
相关产品推荐

