如何使用PySpark DataFrame实现数据透视转换?
PySpark DataFrame 转换方案参考
由于你未提供具体的输入/输出DataFrame结构与示例数据,以下给出几种常见转换场景的最优实现方式,你可以根据实际需求调整:
1. 行转列(Pivot)
适用于将某列的离散值转为新列的场景,例如将分类字段的不同取值拆分为独立列。
- 输入示例:
id category value 1 A 10 1 B 20 2 A 15 - 输出示例:
id A B 1 10 20 2 15 null
实现代码:
from pyspark.sql import functions as F # 按id分组,以category为基准转列,聚合逻辑根据需求选择first/sum等 output_df = input_df.groupBy("id").pivot("category").agg(F.first("value"))
2. 列转行(Unpivot)
适用于将多列合并为键值对行的场景,与Pivot操作相反。
- 输入示例:
id A B 1 10 20 2 15 25 - 输出示例:
id category value 1 A 10 1 B 20 2 A 15 2 B 25
实现代码:
from pyspark.sql import functions as F # stack(n, 键1, 值1, 键2, 值2...) 中n为键值对数量 output_df = input_df.select( "id", F.expr("stack(2, 'A', A, 'B', B) as (category, value)") ).filter(F.col("value").isNotNull()) # 过滤空值(可选)
3. 嵌套字段处理
展开Array类型字段
# 将array字段拆分为多行 output_df = input_df.select("id", F.explode("array_column").alias("element")) # 保留原数组索引(可选) output_df = input_df.select("id", F.posexplode("array_column").alias("index", "element"))
展开Struct类型字段
# 直接展开struct的所有子字段 output_df = input_df.select("id", "struct_column.*") # 选择指定子字段并重命名 output_df = input_df.select( "id", F.col("struct_column.field1").alias("new_field1"), F.col("struct_column.field2").alias("new_field2") )
如果你能补充具体的输入/输出Schema和示例数据,我可以给出更精准的最优实现方案。
内容的提问来源于stack exchange,提问作者DEVEN MALI
相关产品推荐
相关产品推荐

