【问题标题】:Spark serialization (?) with traits weirdness and discrepancies between 1.6.0 ang 2.1.1Spark 序列化 (?) 具有 1.6.0 和 2.1.1 之间的特征怪异和差异
【发布时间】:2017-06-21 00:05:27
【问题描述】:

我知道对此有很多问题,但我无法找到系统的解释来说明究竟什么需要可序列化(以及何时序列化)......以及如何验证此要求。

考虑一下:

class Baz
trait Bar { val baz = new Baz; def bar(i: Int) = baz }
case object Foo extends Bar { def foo = sc.parallelize(1 to 1).map(bar).collect }
Foo.foo

这有效,并返回Array(null) 这对任何人都有意义吗???

如果我将 val 更改为 lazy val,那么它停止工作,并抛出 NotSerializableException,这是有道理的 - 它在远程端初始化 baz,然后无法把它退回。 但是为什么在第一种情况下它会愉快地用null替换它???

如果我写它,几乎是我能想到的任何其他方式 - 例如,将 bar 定义从 trait 移动到对象,或者将 bar 调用替换为 _ => baz - 它也停止工作,并抱怨Task is not serializable.

返回一个在 trait 中定义的 val 的方法是什么,这使得它只是将其写为 null 而不是?有什么想法吗?

更新 上述行为发生在带有 spark 2.1.1 的 scala 2.11 上。 Scala 2.10 (spark 1.6.0) 确实抛出异常,抱怨Baz 不可序列化......所以,这似乎是一种回归。

另外我注意到在 spark 1.6.0 上,这样的东西可以正常工作:

   object Foo { def foo = sc.parallelize(1 to 1).map(bar).collect; def bar(i: Int) = i+1 } 
   Foo.foo

但在 spark 2.1.1 上,它抱怨 Foo 不可序列化。这是为什么? 显然,序列化 lambda 还希望序列化 Foo,这 有点 是有道理的......除了它 确实 在 1.6.0 中以某种方式工作,即使我制作 labda实际引用Foo中的其他东西:

   object Foo { 
     var stuff = 10 
     def foo = sc.parallelize(1 to 1).map(bar).collect
     def bar(i: Int) = { stuff += 1; i+1 }
   } 
   Foo.foo
   Foo.stuff

这在 1.6.0 中可以正常工作,但在 2.1.1 中不行。

所以,这里的一个问题是它在 1.6.0 中实际上是如何工作的?我的意思是,Foo 不可序列化,它如何知道另一端的stuff 的值?

另一个明显的问题是 - 为什么它在 2.1.1 中停止工作? 1.6.0 的行为是否存在微妙的问题,我们不应该依赖它吗? 还是只是 2.1.1 的一个 bug?

【问题讨论】:

    标签: scala apache-spark serialization


    【解决方案1】:

    从某种意义上说,这可能不是一个直截了当的答案,它可以为您提供 Spark 中序列化问题背后的确切原因,以及为什么您的用例可能会以这种或另一种方式工作,但是……让我阐明一下这个。

    SparkContext 是一切发生的地方(或者至少是它开始的地方)。在方法中你可以找到clean方法:

    private[spark] def clean[F <: AnyRef](f: F, checkSerializable: Boolean = true): F = {
      ClosureCleaner.clean(f, checkSerializable)
      f
    }
    

    引用它的 scaladoc 你应该得到关于 Spark 如何进行序列化验证的足够信息:

    清理闭包,使其准备好序列化并发送到任务(删除 $outer's 中未引用的变量,更新 REPL 变量)

    如果设置了checkSerializableclean 也会主动检查f 是否可序列化,如果不是,则抛出SparkException

    搜索所有使用该方法的地方可能会帮助您了解您的代码为什么可以工作或不可以工作。就像将clean 方法应用于您的代码一样简单。

    您也可以从RDD.map 运算符开始,您可以在其中找到clean 方法:

    val cleanF = sc.clean(f)
    

    这样,您可能会了解为什么您的代码会根据您使用vallazy val 给出不同的结果。

    我认为最终你的代码可以重写如下:

    // run spark-shell -c spark.driver.allowMultipleContexts=true
    // use :paste -raw
    package org.apache.spark
    
    class Baz
    trait Bar { lazy val baz = new Baz; def bar(i: Int) = baz }
    case object Foo extends Bar {
      val sc = new SparkContext("local[*]", "Clean", new SparkConf)
      def foo = sc.clean(bar _)
    }
    
    // org.apache.spark.Foo.foo
    

    我为此使用了spark-shell,看来代码在今天构建的 Spark 2.3.0-SNAPSHOT 中使用和不使用 lazy 关键字都可以正常工作。

    【讨论】:

    • 实际上,我了解的唯一一件事是vallazy val case 之间的结果差异:) 其他一切似乎都是一大堆......神秘:)
    • 另外,忘了提一下:无论有没有lazy,您编写的代码的工作方式都相同,但它与我在(第一个)sn-p 中所做的事情不同。如果您将foo 更改为sc.parallelize(1 to 1).map(sc.clean(bar)).collect,那么它将以惰性(这是可以理解的)但“工作”的方式抛出,并返回一个带有 null 的数组而没有惰性,这只是不会使任何意义。
    猜你喜欢
    • 2020-10-02
    • 2019-04-16
    • 1970-01-01
    • 2020-09-18
    • 2021-07-04
    • 2019-07-07
    • 2013-03-25
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多