【问题标题】:Serialization of an object used in foreachRDD() when CheckPointingCheckPointing 时 foreachRDD() 中使用的对象的序列化
【发布时间】:2020-04-28 12:25:22
【问题描述】:

根据this question 和我读过的文档,Spark Streaming 的 foreachRDD(someFunction) 将仅在驱动程序进程中执行 someFunction 本身,但如果有在 RDD 上完成的操作然后将在 executor 上完成 - RDD 所在的位置。

以上所有内容也适用于我,尽管我注意到如果我打开检查点,那么似乎 spark 正在尝试序列化 foreachRDD(someFunction) 中的所有内容并发送到某个地方 - 这是对我造成问题,因为使用的对象之一是不可序列化的(即 schemaRegistryClient)。我尝试了 Kryo 序列化程序,但也没有运气。

如果我关闭检查点,序列化问题就会消失。

有没有办法让 Spark 不序列化 foreachRDD(someFunc) 中使用的内容,同时继续使用检查点?

非常感谢。

【问题讨论】:

    标签: apache-spark spark-streaming rdd avro kryo


    【解决方案1】:

    有没有办法让 Spark 不序列化 foreachRDD(someFunc) 中使用的内容,同时继续使用检查点?

    检查点应该与您的问题无关。根本问题是您有一个不可序列化的对象实例,需要将其发送给您的工作人员。

    当你有这样的依赖时,Spark 中有一个通用的模式可以使用。您创建一个带有惰性瞬态属性的object,该属性将在需要时加载到工作节点中:

    object RegisteryWrapper {
      @transient lazy val schemaClient: SchemaRegisteryClient = new SchemaRegisteryClient()
    }
    

    而当你需要在foreachRDD内部使用它时:

    someStream.foreachRDD { 
       rdd => rdd.foreachPartition { iterator => 
           val schemaClient = RegisteryWrapper.schemaClient
           iterator.foreach(schemaClient.send(_))
      }
    }
    

    【讨论】:

    • 非常感谢 Yuval,您的建议有效。你的一个后续问题-关于“具有惰性瞬态属性,在需要时将在工作节点内加载”的声明-根据此 [链接] (artima.com/pins1ed/annotations.html),它指出“Scala 为字段提供了@transient 注释根本不应该序列化” - 因此,如果工作人员需要 schemaRegistryClient,如何在没有序列化的情况下将其发送到工作人员(或除驱动程序之外的任何其他地方)?
    • @howard 它不会发送给工作人员。但是,不要忘记工作人员拥有操作所需的所有 JAR。这意味着需要一个实例,worker 将加载相关的类并创建一个实例。完全不需要序列化。
    • 我明白了,因为 schemaRegistry 被定义为对象变量,所以它在对象定义中,因此工人可以仅从 jar 本身构造它。类似的规则适用于检查点,它没有被序列化,也没有保存在 HDFS 上,但是驱动程序也有 jar 并且可以从对象定义中构建它 - 所以将它作为对象定义,@transient 注释和惰性前缀都是这种模式所必需的去工作。再次感谢@Yuval!
    • @YuvalItzchakov OP 询问有关foreachRDD 序列化对象的问题。是否可以在foreachRDD 中创建NonSerializable 对象,或者Spark 是否在所有foreach* 函数中强制执行序列化 规则?
    • @CᴴᴀZ 如果你在foreachRDD 函数中分配一个对象,Spark 必须确保它是可序列化的,因为它必须将它发送到 RDD 的每个分区(当然这假设你是致电rdd.foreachrdd.map)。另一方面,在rdd.foreachPartition 内部,如果您在那里进行分配,则分配发生在每个分区本地,因此不需要序列化。
    【解决方案2】:

    这里有几件事很重要:

    1. 您不能在工作人员(即 RDD 内部)上执行的代码中使用此客户端。
    2. 您可以使用临时客户端字段创建对象,并在重新启动作业后重新创建它。可以找到如何完成此操作的示例here.
    3. 同样的原则也适用于广播和累加器变量。
    4. 检查点保留数据、作业元数据和代码逻辑。更改代码后,您的检查点将失效。

    【讨论】:

      【解决方案3】:

      问题可能与检查点数据有关,如果您更改了代码中的任何内容,那么您需要删除旧的检查点元数据。

      【讨论】:

        猜你喜欢
        • 2016-12-12
        • 1970-01-01
        • 1970-01-01
        • 2021-12-08
        • 1970-01-01
        • 1970-01-01
        • 2014-02-10
        • 2014-02-08
        • 2014-10-06
        相关资源
        最近更新 更多