理解这里发生了什么的关键是,虽然 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 也需要成为可序列化class 或object 的一部分
按照建议解决此问题的另一种方法是将val strtoint 更改为UserDefinedFunction,这是一个case class,因此是可序列化的,但是您仍然需要确保strToInt 是可序列化的class 或object。