Spark Scala中结构体数组转换时的异常行为问题求助
解决Spark 2.1.1中嵌套序列的类型推断异常及结构体字段重命名问题
我来帮你梳理这个问题的根源,以及对应的解决办法:
一、为什么会出现[[a,b]]这样的异常输出?
在Spark 2.1.1版本中,当你使用Seq(Seq("a","b"))这种嵌套序列创建DataFrame时,Spark的类型推断机制会把内部固定长度的Seq[String]误判为Tuple2[String, String],最终解析成array<struct<_1:string,_2:string>>类型——这就是你看到输出里data列显示为[[a,b]](数组嵌套结构体)的原因,而非你预期的array<string>。
二、先修正DataFrame的类型(如果需要)
如果你原本期望data列是array<string>类型,有两种可靠的方式来创建正确的DataFrame:
方法1:显式指定Schema
通过提前定义Schema,强制Spark按照你想要的类型解析数据:
import org.apache.spark.sql.types._ import spark.implicits._ // 定义目标Schema val targetSchema = StructType(Seq( StructField("id", IntegerType, nullable = false), StructField("data", ArrayType(StringType), nullable = true) )) // 基于Schema创建DataFrame val df = spark.createDataFrame( spark.sparkContext.parallelize(Seq( (1, Seq("a","b")), (2, Seq("c","d")) )), targetSchema ) df.show(false) df.printSchema()
此时输出会是你预期的:
+---+-------+ |id |data | +---+-------+ |1 |[a,b] | |2 |[c,d] | +---+-------+ root |-- id: integer (nullable = false) |-- data: array (nullable = true) | |-- element: string (containsNull = true)
方法2:转换已有错误类型的DataFrame
如果已经有了那个错误类型的DataFrame,可以通过UDF把结构体数组转换为字符串数组:
import org.apache.spark.sql.functions._ val flattenStructUdf = udf((arr: Seq[(String, String)]) => arr.flatMap { case (a, b) => Seq(a, b) }) val correctedDf = df.withColumn("data", flattenStructUdf($"data"))
三、重命名结构体中的字段(针对已有的结构体数组)
如果你确实需要保留array<struct>类型,只是想重命名结构体里的_1和_2字段,由于Spark 2.1.1还没有transform这种数组操作函数,我们可以用UDF来实现:
步骤1:定义样例类对应重命名后的结构体
case class NamedElement(first: String, second: String)
步骤2:编写UDF转换数组中的每个结构体
val renameStructUdf = udf((arr: Seq[(String, String)]) => { arr.map { case (val1, val2) => NamedElement(val1, val2) } }) // 应用UDF重命名字段 val renamedDf = df.withColumn("data", renameStructUdf($"data")) renamedDf.printSchema()
此时的Schema会变成:
root |-- id: integer (nullable = false) |-- data: array (nullable = true) | |-- element: struct (containsNull = true) | | |-- first: string (nullable = true) | | |-- second: string (nullable = true)
总结
- Spark 2.x早期版本的类型推断对固定长度嵌套序列不够智能,显式指定Schema是避免这类问题的最佳实践
- 针对结构体字段重命名,在低版本Spark中可以通过UDF结合样例类实现,高版本(Spark 2.4+)则可以用
transform函数更简洁地处理
内容的提问来源于stack exchange,提问作者Raphael Roth
相关产品推荐
相关产品推荐

