Spark:如何将含多列的单行拆分为多行(非RDD实现)
实现DataFrame列转行(Unpivot)的高效方法
当然有比RDD更简洁高效的DataFrame API方案来解决这个需求!你要做的是**列转行(Unpivot)**操作,下面给你两种实用的实现方式,都不需要用到RDD:
方法1:使用array + explode(最直观)
这种方法适合只需要提取列值的场景,步骤简单直接:把要转换的colA、colB、colC打包成一个数组,再用explode函数把数组拆分成多行,同时保留ID和String列:
import org.apache.spark.sql.functions.{col, explode, array} val unpivotedDF = df .withColumn("resultColumn", explode(array(col("colA"), col("colB"), col("colC")))) .select("ID", "String", "resultColumn")
运行后就能得到你期望的结构:每一行对应原DataFrame中一行的一个列值,ID和String会自动重复匹配。
方法2:使用stack函数(支持保留原列名)
如果你之后需要知道每个值来自原表的哪一列,可以用Spark SQL的stack函数,它能同时生成原列名和对应值,之后按需取舍即可:
import org.apache.spark.sql.functions._ val unpivotedDF = df.selectExpr( "ID", "String", "stack(3, 'colA', colA, 'colB', colB, 'colC', colC) as (original_column, resultColumn)" ) // 如果不需要原列名,直接drop掉就行 // .drop("original_column")
stack的第一个参数是要转换的列数量,后面跟着成对的(列名字符串,列值),它会自动把这三列拆分成3行。
为什么不用RDD?
这两种方法都是基于Spark的Catalyst优化器的DataFrame操作,比RDD的手动转换更高效——Spark会自动做 predicate pushdown、列裁剪等优化,代码也更简洁易维护。
内容的提问来源于stack exchange,提问作者Pavel Orlov
相关产品推荐
相关产品推荐

