【发布时间】:2021-03-15 05:27:17
【问题描述】:
我使用 Spark 2.4 已经有一段时间了,最近几天才开始切换到 Spark 3.0。切换到 Spark 3.0 运行 udf((x: Int) => x, IntegerType) 后出现此错误:
Caused by: org.apache.spark.sql.AnalysisException: 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;
解决方案是由 Spark 自己提出的,经过一段时间的谷歌搜索后,我进入了 Spark 迁移指南页面:
在 Spark 3.0 中,默认情况下不允许使用 org.apache.spark.sql.functions.udf(AnyRef, DataType)。建议删除返回类型参数以自动切换到类型化 Scala udf,或者将 spark.sql.legacy.allowUntypedScalaUDF 设置为 true 以继续使用它。在 Spark 版本 2.4 及更低版本中,如果 org.apache.spark.sql.functions.udf(AnyRef, DataType) 获得带有原始类型参数的 Scala 闭包,则如果输入值为 null,则返回的 UDF 返回 null。但是,在 Spark 3.0 中,如果输入值为 null,UDF 将返回 Java 类型的默认值。例如,val f = udf((x: Int) => x, IntegerType), f($"x") 在 Spark 2.4 及更低版本中如果 x 为 null,则返回 null,在 Spark 3.0 中返回 0。引入此行为更改是因为 Spark 3.0 默认使用 Scala 2.12 构建。
我注意到我通常使用function.udf API 的方式是udf(AnyRef, DataType),称为UnTyped Scala UDF,而建议的解决方案是udf(AnyRef),称为Typed Scala UDF。
- 据我了解,第一个看起来比第二个更严格的类型,其中第一个明确定义了其输出类型,而第二个没有,因此我对为什么将其称为 UnTyped 感到困惑。
- 函数也被传递给
udf,即(x:Int) => x,显然定义了它的输入类型,但Spark声称You're using untyped Scala UDF, which does not have the input type information?
我的理解正确吗?即使经过更深入的搜索,我仍然找不到任何解释什么是 UnTyped Scala UDF 和什么是 Typed Scala UDF 的材料。
所以我的问题是:它们是什么?它们有什么区别?
【问题讨论】:
标签: scala apache-spark user-defined-functions