将DataFrame条目转为含Any类型成员的Case Class遇问题求助
Spark UDF无法处理Any类型:自定义编码器解决方案
我有一个包含不同数据类型列的DataFrame,结构如下:
+-------+-------+-------+ |column1|column2|column3| +-------+-------+-------+ | 1| a| 0.1| | 2| b| 0.2| | 3| c| 0.3| +-------+-------+-------+
其中column1是Int类型,column2是String类型,column3是Float类型。
我尝试通过UDF把所有列的条目转换成下面的Case Class:
case class Annotation(lastUpdate: String, value: Any)
执行了以下代码:
val columns = df.columns val myUDF= udf { in: Any => Annotation("dummy", in) } val finalDF = columns.foldLeft(df){ (tempDF, colName) => tempDF.withColumn(colName, myUDF(col(colName))) }
但运行时抛出了错误:
Exception in thread "main" java.lang.UnsupportedOperationException: Schema for type scala.Any is not supported at org.apache.spark.sql.catalyst.ScalaReflection$$anonfun$schemaFor$1.apply(ScalaReflection.scala:762) at org.apache.spark.sql.catalyst.ScalaReflection$$anonfun$schemaFor$1.apply(ScalaReflection.scala:704) at scala.reflect.internal.tpe.TypeConstraints$UndoLog.undo(TypeConstraints.scala:56) at org.apache.spark.sql.catalyst.ScalaReflection$class.cleanUpReflectionObjects(ScalaReflection.scala:809) at org.apache.spark.sql.catalyst.ScalaReflection$.cleanUpReflectionObjects(ScalaReflection.scala:39) at org.apache.spark.sql.catalyst.ScalaReflection$.schemaFor(ScalaReflection.scala:703) at org.apache.spark.sql.catalyst.ScalaReflection$$anonfun$schemaFor$1$$anonfun$apply$6.apply(ScalaReflection.scala:758) at org.apache.spark.sql.catalyst.ScalaReflection$$anonfun$schemaFor$1$$anonfun$apply$6.apply(ScalaReflection.scala:757) at scala.collection.TraversableLike$$anonfun$map$1.apply(TraversableLike.scala:234) at scala.collection.TraversableLike$$anonfun$map$1.apply(TraversableLike.scala:234) at scala.collection.immutable.List.foreach(List.scala:381)
我知道可以通过自定义编码器解决,但不清楚具体怎么应用,希望得到可行方案。
问题根源
Spark的SQL引擎依赖明确的Schema来处理数据,而scala.Any是一个无具体类型信息的父类,Spark无法推断出它对应的SQL数据类型,所以直接用Any作为Case Class的字段类型会抛出Schema不支持的错误。
解决方案1:泛型Case Class + 自定义Encoder(推荐)
这个方案能保留原始列的类型信息,后续还能对value字段进行类型相关的操作。
步骤1:定义泛型Case Class
把Annotation改成泛型类,让value字段保留具体的类型:
case class Annotation[T](lastUpdate: String, value: T)
步骤2:编写自定义Encoder
我们需要为Annotation[T]实现Encoder,让Spark能正确解析它的Schema:
import org.apache.spark.sql.catalyst.encoders.ExpressionEncoder import org.apache.spark.sql.catalyst.expressions.{CreateNamedStruct, Literal} import org.apache.spark.sql.types.{StringType, StructField, StructType} import org.apache.spark.sql.Encoder def annotationEncoder[T](implicit encoder: Encoder[T]): Encoder[Annotation[T]] = { // 利用传入的T类型的编码器获取value的表达式 val valueExpr = encoder.toRowExpression // 构造Struct类型的表达式,包含lastUpdate和value两个字段 val structExpr = CreateNamedStruct( Literal("lastUpdate"), Literal("dummy").expr, Literal("value"), valueExpr ) // 定义对应的Schema val schema = StructType(Seq( StructField("lastUpdate", StringType), StructField("value", encoder.schema) )) // 创建并返回ExpressionEncoder ExpressionEncoder[Annotation[T]](structExpr, schema) }
步骤3:创建泛型UDF并应用
我们可以为不同类型的列生成对应的UDF,或者通过类型匹配动态处理所有列:
import org.apache.spark.sql.functions.udf import org.apache.spark.sql.types._ // 生成泛型UDF的方法 def createAnnotationUDF[T](implicit encoder: Encoder[T]) = { udf((value: T) => Annotation("dummy", value))(annotationEncoder[T]) } // 使用foldLeft遍历所有列,根据列类型匹配对应的UDF val finalDF = columns.foldLeft(df) { (tempDF, colName) => val colType = tempDF.schema(colName).dataType val annotationUdf = colType match { case IntegerType => createAnnotationUDF[Int] case StringType => createAnnotationUDF[String] case FloatType => createAnnotationUDF[Float] // 可以根据需要扩展支持更多类型,比如DoubleType、LongType等 case _ => throw new IllegalArgumentException(s"Unsupported column type: $colType") } tempDF.withColumn(colName, annotationUdf(col(colName))) }
解决方案2:将Value转为String(简单但受限)
如果后续不需要对value进行数值运算,只是想存储值,可以把value转换成String类型,这样Spark就能直接推断Schema:
步骤1:修改Case Class
case class Annotation(lastUpdate: String, value: String)
步骤2:调整UDF并应用
val myUDF = udf { in: Any => Annotation("dummy", in.toString) } val finalDF = columns.foldLeft(df){ (tempDF, colName) => tempDF.withColumn(colName, myUDF(col(colName))) }
这个方法实现简单,但缺点是value变成了字符串,后续需要做数值操作时得再转换回来。
内容的提问来源于stack exchange,提问作者mongolol
相关产品推荐
相关产品推荐

