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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.18 00:00:13