You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.28 14:15:55