【发布时间】:2018-05-18 08:53:55
【问题描述】:
【问题讨论】:
标签: java scala performance apache-spark jvm
【问题讨论】:
标签: java scala performance apache-spark jvm
正如已经说过的广播变量是一回事。
另一个是并发问题。看看这段代码:
var counter = 0
var rdd = sc.parallelize(data)
rdd.foreach(x => counter += x)
println(counter)
结果可能会有所不同,具体取决于是在本地执行还是在部署在集群上的 Spark(使用不同的 JVM)上执行。在后一种情况下,parallelize 方法在执行程序之间拆分计算。计算闭包(每个节点执行任务所需的环境),这意味着每个执行程序都会收到counter 的副本。每个执行器都看到自己的变量副本,因此计算结果为 0,因为没有一个执行器引用了正确的对象。另一方面,在一个 JVM 中counter 对每个工作人员都是可见的。
当然有一种方法可以避免这种情况 - 使用 Acumulators (see here)。
最后但同样重要的是,当在内存中持久化RDDs 时(默认cache 方法存储级别为MEMORY_ONLY),它将在单个JVM 中可见。这也可以通过使用OFF_HEAP 来克服(这在 2.4.0 中是实验性的)。更多here。
【讨论】:
最大的可能优势是共享内存,特别是处理广播对象。因为这些对象被认为是只读的,所以可以在多个线程之间共享。
在使用单个任务/执行程序的场景中,每个 JVM 都需要一个副本,因此 N 个任务有 N 个副本。对于大型对象,这可能是一个严重的开销。
同样的逻辑可以应用于其他共享对象。
【讨论】: