Spark UDF返回Row报错:向Map添加Struct元素失败排查
问题:Spark中向Map类型字段添加Struct条目报错处理
业务场景与环境
使用Scala 2.13.13 + Spark 3.3.1,需要向DataFrame中Map类型的字段(String到Struct的映射)添加新的Struct条目,原始代码如下:
val json = """ [ { "info" : { "1234" : { "name": "John Smith", "age": 29, "gender": "male" } } } ] """ val personSchema = new StructType() .add("name", StringType) .add("age", IntegerType) .add("gender", StringType) val schema = new StructType().add("info", MapType(StringType, personSchema)) val spark = SparkSession.builder() .master("local[*]") .getOrCreate() spark.sparkContext.setLogLevel("ERROR") import spark.implicits._ val df = spark.read.schema(schema).json(Seq(json).toDS) df.show(false)
预期结果
希望动态生成一个Person条目添加到info的Map中,最终JSON结构如下:
{ "info": { "1234": { "name": "John Smith", "age": 29, "gender": "male" }, "456": { "name": "Robert Jones", "age": 35, "gender": "male" } } }
尝试的UDF及报错
编写了以下UDF实现,但运行时报错:
val addPersonUDF = udf((infoMap: Map[String, Row]) => { infoMap + ("456" -> new GenericRowWithSchema(Array("Robert Jones", 35, "male"), personSchema)) }) df.select(col("*"), addPersonUDF(col("info"))).show(false)
报错信息:
Exception in thread "main" java.lang.UnsupportedOperationException: Schema for type org.apache.spark.sql.Row is not supported at org.apache.spark.sql.errors.QueryExecutionErrors$.schemaForTypeUnsupportedError(QueryExecutionErrors.scala:1193) at org.apache.spark.sql.catalyst.ScalaReflection$.$anonfun$schemaFor$1(ScalaReflection.scala:802) at scala.reflect.internal.tpe.TypeConstraints$UndoLog.undo(TypeConstraints.scala:73) at org.apache.spark.sql.catalyst.ScalaReflection.cleanUpReflectionObjects(ScalaReflection.scala:948) at org.apache.spark.sql.catalyst.ScalaReflection.cleanUpReflectionObjects$(ScalaReflection.scala:947) at org.apache.spark.sql.catalyst.ScalaReflection$.cleanUpReflectionObjects(ScalaReflection.scala:51) at org.apache.spark.sql.catalyst.ScalaReflection$.schemaFor(ScalaReflection.scala:718) at org.apache.spark.sql.catalyst.ScalaReflection$.$anonfun$schemaFor$1(ScalaReflection.scala:744) at scala.reflect.internal.tpe.TypeConstraints$UndoLog.undo(TypeConstraints.scala:73) at org.apache.spark.sql.catalyst.ScalaReflection.cleanUpReflectionObjects(ScalaReflection.scala:948) at org.apache.spark.sql.catalyst.ScalaReflection.cleanUpReflectionObjects$(ScalaReflection.scala:947) at org.apache.spark.sql.catalyst.ScalaReflection$.cleanUpReflectionObjects(ScalaReflection.scala:51) at org.apache.spark.sql.catalyst.ScalaReflection$.schemaFor(ScalaReflection.scala:718) at org.apache.spark.sql.catalyst.ScalaReflection$.schemaFor(ScalaReflection.scala:714) at org.apache.spark.sql.functions$.$anonfun$udf$6(functions.scala:5124) at scala.Option.getOrElse(Option.scala:201) at org.apache.spark.sql.functions$.udf(functions.scala:5124)
问题原因与解决方法
问题原因
Spark自动推导UDF的返回Schema时,无法识别Row作为Map的值类型。Row是Spark内部通用行对象,反射机制无法直接为其生成对应的StructSchema,因此抛出Schema不支持的错误。
方法一:使用Spark内置函数(推荐)
无需自定义UDF,直接用map_concat结合lit和struct函数实现,更贴合Spark的优化逻辑:
import org.apache.spark.sql.functions.{col, lit, map_concat, struct} // 构造新的Map条目,key为"456",value为对应的Struct val newEntry = lit(Map("456" -> struct( lit("Robert Jones").as("name"), lit(35).as("age"), lit("male").as("gender") ))) // 合并原Map和新条目 val resultDf = df.withColumn("info", map_concat(col("info"), newEntry)) resultDf.show(false) resultDf.toJSON.show(false)
方法二:显式指定UDF的返回Schema
如果一定要用UDF,需要手动指定返回的Schema,而不是依赖自动推导:
import org.apache.spark.sql.types.MapType // 显式指定UDF返回的Schema为原Map的Schema val addPersonUDF = udf( (infoMap: Map[String, Row]) => { infoMap + ("456" -> new GenericRowWithSchema(Array("Robert Jones", 35, "male"), personSchema)) }, MapType(StringType, personSchema) // 手动指定返回Schema ) val resultDf = df.withColumn("info", addPersonUDF(col("info"))) resultDf.show(false)
内容的提问来源于stack exchange,提问作者Pepria
相关产品推荐
相关产品推荐

