【问题标题】:Why use a singleton to wrap broadcast variables?为什么使用单例来包装广播变量?
【发布时间】:2018-06-13 21:04:52
【问题描述】:

我在https://github.com/apache/spark/blob/master/examples/src/main/java/org/apache/spark/examples/streaming/JavaRecoverableNetworkWordCount.java看到了一个示例代码

代码使用单例来包装广播变量,如下所示:

class JavaWordBlacklist {

private static volatile Broadcast<List<String>> instance = null;

public static Broadcast<List<String>> getInstance(JavaSparkContext jsc) {
  if (instance == null) {
    synchronized (JavaWordBlacklist.class) {
      if (instance == null) {
        List<String> wordBlacklist = Arrays.asList("a", "b", "c");
        instance = jsc.broadcast(wordBlacklist);
      }
    }
  }
  return instance;
}
}

并在wordCounts.foreachRDD((rdd, time) -&gt; {...}中初始化广播变量

我的问题是为什么不在父类中声明private static volatile Broadcast&lt;List&lt;String&gt;&gt; instance = null;,即JavaRecoverableNetworkWordCount

(在我看来,由于广播变量在foreachRDD()中初始化,在单个驱动线程中执行,这里不会发生竞争条件,所以单例保护是不必要的。)

【问题讨论】:

    标签: java apache-spark spark-streaming


    【解决方案1】:

    这样做是为了解决检查点恢复中出现的问题。请记住,检查点仅捕获元数据和/或分布式状态,而不是广播变量、累加器和本地对象。应用程序从检查点重新启动后,必须手动恢复所有状态。

    不解决你的问题:

    由于广播变量在foreachRDD()中初始化,在单个驱动线程中执行,

    驱动程序不是单线程的,访问广播变量的目的与数据处理(簿记、报告)不同。也可以同时被多个流访问。

    【讨论】:

      猜你喜欢
      • 2016-04-10
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多