【问题标题】:Spark Broadcast variable causing failure when reloading state from a checkpoint directory after a reboot重启后从检查点目录重新加载状态时,Spark Broadcast 变量导致失败
【发布时间】:2018-12-07 11:54:20
【问题描述】:

我有一个从 Kafka 读取数据的 spark 流应用程序(使用 spark 1.6.1)。我正在使用 hdfs 上的 spark 检查点目录从故障中恢复。

代码在每批开始时使用带有映射的广播变量,如下所示

public static void execute(JavaPairDStream<String, MyEvent> events) {
  final Broadcast<Map<String, String>> appConfig =  MyStreamConsumer.ApplicationConfig.getInstance(new    JavaSparkContext(events.context().sparkContext()));

当我在运行时提交作业时,一切正常,来自 kafka 的所有事件都得到正确处理。从故障中恢复时会出现问题(通过重新启动运行 spark 的机器进行测试) - spark 流应用程序实际上正确启动,并且在 UI 中一切看起来都很好,作业正在运行,但是一旦通过以下异常发送数据调用此方法时遇到(&作业崩溃):

appConfig.value() (the broadcast variable from the start!)

由于 Spark 作业错误而失败

  Caused by: java.lang.ClassCastException: org.apache.spark.util.SerializableConfiguration cannot be cast to java.util.Map

如果我在 spark UI 中终止驱动程序并从命令行重新提交作业,一切都会再次正常运行。 但它对我们产品的要求是它可以从故障中自动恢复,甚至只是重新启动任何集群节点,所以我必须修复上述问题。 这个问题肯定与使用 Broadcast 变量和在重启后从 spark checkpoint 目录加载状态有关

另外请注意,我确实正确地创建了广播实例(懒惰/单例):

public static Broadcast<Map<String, String>> getInstance(JavaSparkContext sparkContext) {
    if (instance == null) {

我确实意识到这个问题似乎与: Is it possible to recover an broadcast value from Spark-streaming checkpoint

但我无法按照他们的说明解决问题

【问题讨论】:

    标签: apache-spark


    【解决方案1】:

    回答我自己的问题,因为其他人可能会觉得它很有用: 奇怪但以下似乎已经解决了这个问题,如果我将以下代码移动到执行方法的开头: appConfig.value() 并将其分配给那里的法线贴图变量。

    然后,如果我只是在我的匿名 FlatMapFunction 代码中使用这个 map 变量而不是 appConfig.value() - 即使在重新启动后一切正常。

    同样,不知道为什么会这样,但确实如此......

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2019-08-07
      • 2015-06-27
      • 1970-01-01
      • 2021-05-07
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多