如何高效转置含单行多整数列的PySpark DataFrame
高效转置Spark DataFrame的方案
原始DataFrame结构:
列名:A、B、C、D、E、F 列类型:全部为Long 数据行:[4, 2, 7, 9, 1, 3]目标转置后结构:
列名:name、value 列类型:String、Long 数据行: A 4 B 2 C 7 D 9 E 1 F 3当前使用
df.collect()实现,耗时过长,需要高效替代方案。
核心问题说明
df.collect()会把分布式存储的DataFrame数据全部拉到Driver节点,数据量较大时必然卡顿甚至触发内存溢出,必须用Spark原生的分布式操作来实现转置。
方法1:使用stack内置函数(通用所有Spark版本)
stack是Spark SQL原生函数,专门用于将多列转为多行,全程分布式执行,性能远优于collect()。
Scala代码示例
import org.apache.spark.sql.functions._ // 构造示例DataFrame val df = spark.createDataFrame(Seq((4L, 2L, 7L, 9L, 1L, 3L))) .toDF("A", "B", "C", "D", "E", "F") // 执行转置 val transposedDF = df.selectExpr( "stack(6, 'A', A, 'B', B, 'C', C, 'D', D, 'E', E, 'F', F) as (name, value)" ) transposedDF.show()
Python代码示例
from pyspark.sql import functions as F // 构造示例DataFrame df = spark.createDataFrame([(4, 2, 7, 9, 1, 3)], ["A", "B", "C", "D", "E", "F"]) // 执行转置 transposed_df = df.selectExpr( "stack(6, 'A', A, 'B', B, 'C', C, 'D', D, 'E', E, 'F', F) as (name, value)" ) transposed_df.show()
- 参数说明:
stack(n, col1_name, col1_val, ...)中,n是要转置的列总数,后续依次传入列名的字符串和对应的列对象。
方法2:使用melt函数(Spark 3.1+版本可用)
Spark 3.1及以上版本支持Pandas风格的melt函数,代码更简洁,语义更清晰。
Scala代码示例
import org.apache.spark.sql.functions._ import spark.implicits._ val df = spark.createDataFrame(Seq((4L, 2L, 7L, 9L, 1L, 3L))) .toDF("A", "B", "C", "D", "E", "F") val transposedDF = df.melt( id_vars = Seq.empty[String], // 无需要保留的标识列 value_vars = df.columns, // 指定要转置的所有列 var_name = "name", // 新列"name"对应原列名 value_name = "value" // 新列"value"对应原列值 ) transposedDF.show()
Python代码示例
from pyspark.sql import functions as F df = spark.createDataFrame([(4, 2, 7, 9, 1, 3)], ["A", "B", "C", "D", "E", "F"]) transposed_df = df.melt( id_vars=[], value_vars=df.columns, var_name="name", value_name="value" ) transposed_df.show()
内容的提问来源于stack exchange,提问作者DQd
相关产品推荐
相关产品推荐

