Spark 3.5+中修改返回StructType的已弃用UDF方法
问题描述
我们正在将GCP Dataproc集群升级至2.2debian12镜像,对应Spark版本3.5.0、Scala版本2.12.18。升级后遇到UDF兼容性问题:自Spark 3.0.0起,带返回类型参数的Scala UDF方法已被弃用,我们有部分UDF需要返回IntegerType和StructType。此前在Spark 3.4、Scala 2.12.16的Dataproc 2.0镜像中,可通过设置属性--properties=spark.sql.legacy.allowUntypedScalaUDF=true正常运行作业,但迁移至新版本后出现如下错误:
AnalysisException: [UNTYPED_SCALA_UDF] You're using untyped Scala UDF, which does not have the input type information. Spark may blindly pass null to the Scala closure with primitive-type argument, and the closure will see the default value of the Java type for the null argument, e.g. udf((x: Int) => x, IntegerType), the result is 0 for null input. To get rid of this error, you could: 1. use typed Scala UDF APIs(without return type parameter), e.g. udf((x: Int) => x). 2. use Java UDF APIs, e.g. udf(new UDF1[String, Integer] { override def call(s: String): Integer = s.length() }, IntegerType), if input types are all non primitive. 3. set "spark.sql.legacy.allowUntypedScalaUDF" to "true" and use this API with caution.
修正后的返回StructType的UDF代码
针对返回StructType的UDF,采用Spark类型安全的UDF API,通过定义case class映射Struct结构,让Spark自动推导返回类型,避免手动指定返回类型参数。
步骤1:定义对应Struct的Case Class
首先定义与目标StructType字段一一对应的case class,nullable字段用Option[T]包装:
case class TimestampResult( t: java.sql.Timestamp, tl: Option[java.sql.Timestamp], tlon: Option[java.sql.Timestamp], o: Option[Int] )
步骤2:修改UDF实现
将UDF的返回值从Row改为TimestampResult实例,同时移除手动指定的返回类型参数:
// 先定义对应Struct的case class case class TimestampResult( t: java.sql.Timestamp, tl: Option[java.sql.Timestamp], tlon: Option[java.sql.Timestamp], o: Option[Int] ) val timestampParserLocal: UserDefinedFunction = udf((in: Any, pattern: String, defaultTimezone: String, dropSubSeconds: Boolean, optionalTZ: Boolean, fixInferTimestampsBackwardsCompatibility: Boolean) => { val (t, tl, tb, offset) = in match { case ts: java.sql.Timestamp => if (fixInferTimestampsBackwardsCompatibility) { val format1 = new java.text.SimpleDateFormat(pattern) format1.setLenient(false) val value = format1.format(ts) parseWithLocal(value, pattern, defaultTimezone) } else { // 为nullable字段包装Option,非null值用Some,null用None (ts, Some(ts), Some(getAsLondon(ts)), Some(0)) } case value: String => parseWithLocal(adjTsString(value, dropSubSeconds, optionalTZ), pattern, defaultTimezone) case value: Long => parseWithLocal(value.toString, pattern, defaultTimezone) case _ => throw new RuntimeException(s"Unexpected datatype, in=${in}, class=${in.getClass.getName}") } // 返回case class实例,Spark自动推导StructType TimestampResult(t, tl, tb, offset) })
关键修改点
- 用case class替代手动构造的
Row,Spark会自动将case class映射为对应的StructType,无需手动定义StructType参数。 - 原StructType中
nullable=true的字段,对应case class中的Option[T]类型,避免null值引发的类型安全问题。 - 确保
parseWithLocal方法返回的元组中,对应nullable字段的值是Option[T]类型(若原逻辑返回null,需改为None;非null值用Some包装)。
内容的提问来源于stack exchange,提问作者chanchal ahuja
相关产品推荐
相关产品推荐

