【问题标题】:What are Untyped Scala UDF and Typed Scala UDF? What are their differences?什么是非类型化 Scala UDF 和类型化 Scala UDF?他们有什么区别?
【发布时间】: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 构建。

来源:Spark Migration Guide

我注意到我通常使用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


    【解决方案1】:

    这并没有回答您关于不同 UDF 是什么的原始问题,但如果您想摆脱错误,在 Python 中,您可以在脚本中包含以下行:spark.sql("set spark.sql.legacy.allowUntypedScalaUDF=true")。

    【讨论】:

      【解决方案2】:

      在类型化的 Scala UDF 中,UDF 知道作为参数传递的列的类型,而在非类型化的 Scala UDF 中,UDF 不知道作为参数传递的列的类型

      在创建类型化 scala UDF 时,作为参数传递的列类型和 UDF 的输出是从函数参数和输出类型推断出来的,而在创建非类型化 scala UDF 时,根本没有类型推断,无论是参数还是输出.

      令人困惑的是,在创建类型化 UDF 时,类型是从函数推断出来的,而不是作为参数显式传递的。更明确地说,您可以编写类型化 UDF 创建如下:

      val my_typed_udf = udf[Int, Int]((x: Int) => Int)
      

      现在,让我们看看你提出的两点。

      据我了解,第一个(例如udf(AnyRef, DataType))看起来比第二个(例如udf(AnyRef))更严格,其中第一个明确定义了其输出类型,而第二个没有,因此我感到困惑为什么它被称为 UnTyped。

      根据spark functions scaladoc,将函数转换为UDF的udf函数的签名实际上是第一个:

      def udf(f: AnyRef, dataType: DataType): UserDefinedFunction 
      

      第二个:

      def udf[RT: TypeTag, A1: TypeTag](f: Function1[A1, RT]): UserDefinedFunction
      

      所以第二个实际上比第一个有更多类型,因为第二个考虑了作为参数传递的函数的类型,而第一个删除了函数的类型。

      这就是为什么在第一个你需要定义返回类型的原因,因为 spark 需要这个信息但不能从作为参数传递的函数中推断它,因为它的返回类型被删除了,而在第二个中,返回类型是从作为参数传递的函数。

      函数也被传递给udf,即(x:Int) => x,显然定义了它的输入类型但Spark声称You're using untyped Scala UDF, which does not have the input type information?

      这里重要的不是函数,而是 Spark 如何从该函数创建 UDF。

      在这两种情况下,要转换为 UDF 的函数都定义了其输入和返回类型,但这些类型会被删除,并且在使用 udf(AnyRef, DataType) 创建 UDF 时不会考虑在内。

      【讨论】:

        猜你喜欢
        • 2012-06-03
        • 2012-05-02
        • 2012-02-02
        • 2011-01-19
        • 2022-01-22
        • 2010-10-02
        • 2012-04-14
        相关资源
        最近更新 更多