如何将Spark中的WrappedArray转为DataFrame列?求更优方案
问题描述
我有如下结构的Avro文件,希望将其中的WrappedArray值转换为特定列:
root |-- reportId: string (nullable = true) |-- row: string (nullable = true) |-- col: struct (nullable = true) | |-- fieldRec: array (nullable = true) | | |-- element: struct (containsNull = true) | | | |-- value: string (nullable = true)
读取并选择指定列的代码:
//Reading and selecting specific column from Avro val empDF = spark.read.format("avro").load("/PRD/DRV/EMP_DETAILS/fact_date=20220404") .select("col.fieldRec.value") //I have selected now the "value" which is having WrappedArray empDF.show()
执行后结果:
+------------------------------------+ | value | +------------------------------------+ |WrappedArray(EmployeeId,Name,salary)| |WrappedArray(100,James,300000) | |WrappedArray(101,John,400000) | |WrappedArray(102,Fern,350000) | |WrappedArray(103,Philip,310000) | +------------------------------------+
我当前已实现转换逻辑,但想知道是否有更优方案:
import scala.collection.mutable.ArrayBuffer //Defining Case Class case class Employee(employeeId: String, name: String, salary: String) //Reading and selecting specific column from DataFrame val empDF = spark.read.format("avro").load("/PRD/DRV/EMP_DETAILS/fact_date=20220404") .select("col.fieldRec.value") //Select the First record from the DataFrame val firstRec = empDF.first //Filtering out the first record val empDFMinus = empDF.filter(row => row != firstRec) //Identifying the size of the Array which could be used at the later stage val arrayLength = empDF.select(size(col("value"))).first.get(0).toString.toInt - 1 import spark.implicits._ //Converting WrappedArray to a structure which I would like to represent val empDFCols = empDFMinus.as[ArrayBuffer[String]].map(arr => Employee(arr(0),arr(1),arr(2))) empDFCols.printSchema() empDFCols.show()
输出结果:
root |-- employeeId: string (nullable = true) |-- name: string (nullable = true) |-- salary : string (nullable = true) +----------------------------+ |employeeId | name | salary | +----------------------------+ |100 |James | 300000 | |101 |John | 400000 | |102 |Fern | 350000 | |103 |Philip| 310000 | +----------------------------+
优化方案
方案1:固定列数与列名(已知结构场景)
如果明确知道数组元素对应的列名,直接用Spark原生getItem提取数组元素,无需转换为Dataset,性能更优:
import org.apache.spark.sql.functions._ val empDF = spark.read.format("avro") .load("/PRD/DRV/EMP_DETAILS/fact_date=20220404") .select("col.fieldRec.value") // 跳过表头行,提取数据行的数组元素作为对应列 val resultDF = empDF .filter(row => !row.getAs[Seq[String]]("value").sameElements(Seq("EmployeeId", "Name", "salary"))) .select( col("value").getItem(0).alias("employeeId"), col("value").getItem(1).alias("name"), col("value").getItem(2).alias("salary") ) resultDF.printSchema() resultDF.show()
方案2:动态提取列名(适配表头变化场景)
如果表头可能变动,可一次性获取表头和数据行,避免多次触发计算:
import org.apache.spark.sql.functions._ val empDF = spark.read.format("avro") .load("/PRD/DRV/EMP_DETAILS/fact_date=20220404") .select("col.fieldRec.value") // 一次性收集表头和数据行,减少action触发次数 val (header, dataRows) = empDF.rdd.collect() match { case Array(h, rest @ _*) => (h.getAs[Seq[String]]("value"), rest) } // 将数据行转为DataFrame,用表头动态生成列 val dataDF = spark.createDataFrame(dataRows, empDF.schema) val resultDF = dataDF.select( header.zipWithIndex.map { case (colName, idx) => col("value").getItem(idx).alias(colName) }: _* ) resultDF.printSchema() resultDF.show()
方案3:explode + pivot(复杂动态场景)
如果数组长度不固定或需要更灵活的列转换,可结合explode和pivot实现:
import org.apache.spark.sql.functions._ val empDF = spark.read.format("avro") .load("/PRD/DRV/EMP_DETAILS/fact_date=20220404") .select("col.fieldRec.value") // 添加行号区分表头和数据行 val indexedDF = empDF.withColumn("row_num", monotonically_increasing_id()) // 提取表头列名 val header = indexedDF.filter(col("row_num") === 0) .select(explode(col("value")).alias("col_name")) .collect().map(_.getString(0)) // 将数组转为键值对后pivot为列 val resultDF = indexedDF.filter(col("row_num") > 0) .select( col("row_num"), explode(arrays_zip(col("value"), lit(header))).alias("kv") ) .select( col("row_num"), col("kv._1").alias("value"), col("kv._2").alias("col_name") ) .groupBy("row_num") .pivot("col_name") .agg(first("value")) .drop("row_num") resultDF.printSchema() resultDF.show()
优化优势
- 避免原方案中多次调用
first()触发的重复计算,减少Driver端压力 - 采用Spark原生列操作,更适配分布式计算场景
- 动态列名方案无需硬编码,支持表头结构变动
内容的提问来源于stack exchange,提问作者Vijay B
相关产品推荐
相关产品推荐

