【问题标题】:How to apply user-defined function to column (gives "Task not serializable" when adding a column)?如何将用户定义的函数应用于列(添加列时给出“任务不可序列化”)?
【发布时间】:2017-10-29 15:49:27
【问题描述】:

我必须附加由“strToInt”方法生成的这个列,结果证明它是不可序列化的。

def strToInt(colVal : String) : Int = {
  var str = new Array[String](3)
  str(0) = "icmp"; str(1) = "tcp"; str(2) = "udp"
  var i = 0
  for (i <- 0 to str.length-1) {
    if (str(i) == colVal) { return i }
  }
  throw new IllegalStateException("This never happens")
}
val strtoint = udf(strToInt(_:String)).apply(col("Atr 1"))
val newDF = df.withColumn("newCol", strtoint)

我曾尝试以这种方式将函数放入辅助类中,

object Helper extends Serializable {
    def strToInt ...     
                                    }

但没用。

【问题讨论】:

    标签: scala apache-spark apache-spark-sql


    【解决方案1】:

    将您的代码更改为如下,其中函数执行处于withColumn 级别(而不是在定义 UDF 时)。

    // define a UDF
    val strtoint = udf(strToInt _)
    // use it (aka execute)
    val newDF = df.withColumn("newCol", strtoint(col("Atr 1")))
    

    看似的微小变化会改变你创建的内容以及之后的执行方式。

    您可能已经注意到,udf 创建了一个 Spark SQL 可以理解(可以执行)的用户定义函数:

    udf[RT, A1](f: (A1) ⇒ RT): UserDefinedFunction 将 1 个参数的用户定义函数定义为用户定义函数 (UDF)。

    (为了便于理解,我去掉了隐式参数)

    引用UserDefinedFunction的scaladoc:

    用户定义的函数。要创建一个,请使用函数中的udf 函数。

    我不太同意,但“协议”是先注册一个 UDF,然后才能在查询中执行它,比如 withColumnselect 运算符。


    我还将 strToInt 更改为更符合 Scala 习惯(希望也更容易理解)。

    def strToInt(colVal : String) : Int = {
      val strs = Array("icmp", "tcp", "udp")
      strs.indexOf(colVal)
    }
    

    【讨论】:

    • 非常感谢!那行得通!您能否解释一下与当前语法相比,之前的语法做了什么(和预期的)?
    • 希望 cmets 能给你一些见解。关键是您想要定义立即执行 UDF,这不是您想要的。
    • 它与类型有关。当你这样做时:val strtoint = udf(strToInt _) strtoint 的类型是 UserDefinedFunction,这是可序列化的案例类 当你这样做时:val strtoint = udf(strToInt(_:String)).apply(col("Atr 1")) strtoint 是类型 Function1[String, UserDefinedFunction],基本上是一个匿名函数,它生成一个 @987654336 @。 Function1 不可序列化。 Spark 需要它是可序列化的,以便它可以将其发送到在 DataFrame 上执行分布式操作的每个节点。
    • @kapunga 谢谢。我有一个快速跟进的问题。我想保留“对象助手”并将“str”变量作为对象变量(因此可以从外部分配)。但是,这无法使用语法行执行,给出“无法执行用户定义的函数”。我不明白。
    • @thebluephantom 这可能是一个瓶颈。但是,这取决于我们在谈论什么 UDF。简单的是无害的。
    【解决方案2】:

    理解这里发生了什么的关键是,虽然 Scala 是一种函数式编程语言,但它运行在不支持函数类型的 JVM 上。在运行时,任何分配了“匿名”或“lambda”函数的val 实际上都是具有apply 方法的匿名类的实例。因此,假设您有以下内容:

    object helper {
      val isNegative: (Int => Boolean) = (n: Int) => n < 0
    }
    

    这编译成和这个一样的东西:

    object helper {
      val isNegative: Function1[Int, Boolean] = {
        def apply(n: Int): Boolean = n < 0
      }
    }
    

    isNegative 实际上是一个匿名类实例,扩展了 trait Function1。当您改为这样做时:

    object helper {
      def isNegative(n: Int): Boolean = n < 0
    }
    

    现在isNegative 是对象helper 的一个方法。在处理 Spark 时,如果您要这样做:

    // ds is a Dataset[Int]
    ds.filter(isNegative)
    

    在第一种情况下,Spark 必须序列化分配给isNegative 的匿名类,但由于它不可序列化而失败。在第二种情况下,它必须序列化helper,这确实有效,因为如果object 的所有状态都是可序列化的,它就是可序列化的。

    要将此应用于您的问题,当您这样做时:

    val strtoint = udf(strToInt(_:String)).apply(col("Atr 1"))
    

    在运行时strtoint 是一个具有Funtion1[String, UserDefinedFunction] 特征的匿名类实例,这是一个在被调用时生成UserDefinedFunction 的方法。下划线填上后,和这个是一样的:

    val strtoInt: Function1[String, UserDefinedFunction] = new Function1[String, UserDefinedFunction] = {
      def apply(t1: String) = udf(strToInt(t1 :String)).apply(col("Atr 1"))
    }
    

    要尽量减少代码更改,您只需将 val 更改为 def

    def sti = udf(strToInt(_:String)).apply(col("Atr 1"))
    

    现在sti 是它的封闭类的成员函数,如果它是可序列化的,那么就 Spark 而言你应该是好的。这里要记住的另一件事是strToInt 也需要成为可序列化classobject 的一部分

    按照建议解决此问题的另一种方法是将val strtoint 更改为UserDefinedFunction,这是一个case class,因此是可序列化的,但是您仍然需要确保strToInt 是可序列化的classobject

    【讨论】:

    • 你在strtointsti 中做了什么改变?
    • 它从 val 变为 def。 Scala 处理这些的方式有所不同。对于val,您在运行时将Function1 类型的实例分配给strtoint。此实例是 Spark 尝试序列化的内容,但由于 Function1 不可序列化而失败。当您将defsti 一起使用时,它是一些classobject 的成员方法,并且如果classobject 是可序列化的(例如,如果class 是@987654367 @ 或 object 只有可序列化的成员)然后 Spark 没有问题。
    • 天哪,硬东西。 stackoverflow.com/questions/43592742/… 所以,在这篇文章中,他们主张让 def 只是一个 val - 第一类公民。这与您的答案相比如何?接受的答案对我来说更容易理解。请详细说明。
    【解决方案3】:

    这个问题似乎与我遇到的问题相似(在 Java 中)。 我的 udf 函数使用 Cipher 库来加密某些东西,抛出的异常是:

    Caused by: java.io.NotSerializableException: javax.crypto.Cipher Serialization stack: - object not serializable (class: javax.crypto.Cipher, value: javax.crypto.Cipher@625d02ce)

    我无法将“implements Serializable”添加到 Cipher 类,因为它是 Java 提供的库。

    我从这个链接使用了以下解决方案:spark-how-to-call-udf-over-dataset-in-java

    private static UDF1 toUpper = new UDF1<String, String>() {
        public String call(final String str) throws Exception {
            return str.toUpperCase();
        }
    };
    

    注册UDF,就可以使用callUDF函数了。

    import static org.apache.spark.sql.functions.callUDF;
    import static org.apache.spark.sql.functions.col;
    
    sqlContext.udf().register("toUpper", toUpper, DataTypes.StringType);
    peopleDF.select(col("name"),callUDF("toUpper", col("name"))).show();
    

    而不是调用 str.toUpperCase();我调用了我的 Cipher 实例。

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 2016-11-29
      • 2020-11-11
      • 1970-01-01
      • 2018-11-10
      • 1970-01-01
      • 2021-08-12
      • 1970-01-01
      相关资源
      最近更新 更多