【发布时间】:2015-05-14 18:40:57
【问题描述】:
我在使用广播变量时遇到了一个奇怪的行为。每次我使用广播变量时,内容都会为每个节点复制一次,并且永远不会重复使用。
这是 spark-shell --master local[32] 中的一个示例: (当然,这是无用且愚蠢的代码,但它确实显示了行为)
case class Test(a:String)
val test = Test("123")
val bc = sc.broadcast(test)
// On my 32 core machine, I get 33 copies of Test (expected)
// Yourkit profiler shows 33 instances of my object (32 are unreachable)
sc.parallelize((1 to 100)).map(x => bc.value.a).count
// Doing it again, Test copies are not reused and serialized again (now 65 copies, 64 are unreachable)
sc.parallelize((1 to 100)).map(x => bc.value.a).count
在我的例子中,我广播的变量是几百兆字节,由数百万个小对象(很少的哈希图和向量)组成。
每次我在使用它的 RDD 上运行操作时,都会浪费数 GB 的内存,并且垃圾收集器越来越成为瓶颈!
是设计为每次执行新闭包时重新广播变量还是一个错误,我应该重复使用我的副本?
为什么它们在使用后立即无法访问?
本地模式下的 spark-shell 是特殊的?
注意:我使用的是 spark-1.3.1-hadoop2.6
更新1: 根据这个帖子:http://apache-spark-user-list.1001560.n3.nabble.com/How-to-share-a-NonSerializable-variable-among-tasks-in-the-same-worker-node-td11048.html 单例对象在 Spark 1.2.x+ 上不再工作 所以这种解决方法也行不通:
val bcModel = sc.broadcast(bigModel)
object ModelCache {
@transient lazy private val localModel = { bcModel.value }
def getModel = localModel
}
sc.parallelize((1 to 100)).map(x => ModelCache.getModel.someValue)
更新 2: 我也尝试过重用累加器模式但没有成功:
class VoidAccumulatorParams extends AccumulatorParam[BigModel] {
override def addInPlace(r1: BigModel, r2: BigModel): BigModel= { r1 }
override def zero(initialValue: BigModel): BigModel= { initialValue }
}
val acc = sc.accumulator(bigModel, "bigModel")(new VoidAccumulableParams())
sc.parallelize((1 to 100)).map(x => acc.localValue.someValue)
更新3: 看起来单例对象在使用 spark-submit 而不是 scala shell 运行作业时有效。
【问题讨论】:
标签: apache-spark