【问题标题】:Apache Spark broadcast variables are not reusedApache Spark 广播变量不被重用
【发布时间】: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


    【解决方案1】:

    看看广播测试 (BroadcastSuite.scala) 看看会发生什么。

    当您运行第一个作业时,您的对象会被序列化,切割成块并将这些块发送到执行器(通过 BlockManager 机制)。他们从块中反序列化对象并将其用于处理任务。当他们完成时,他们会丢弃 反序列化的对象,但 BlockManager 会缓存序列化数据的块。

    对于第二个作业,对象不需要被序列化和传输。它只是从缓存中反序列化并使用。


    注意事项:一方面,这并不能帮助您避免过多的 GC。另一件事是,我试图通过使用可变状态 (class Test(var a: String) extends Serializable) 并在运行之间对其进行变异来验证上述理论。令我惊讶的是,第二次运行看到了突变状态!所以我要么完全错了,要么在local 模式下错了。我希望有人能说出哪个。 (如果我明天记得的话,我会尝试自己进一步测试。)

    【讨论】:

    • 感谢您的反馈。我有点惊讶它每次都反序列化对象。对我来说,它看起来像一个功能完成了一半。如果在 context.broadcast(someVar, keepAlive=true) 上有一个标志之类的东西会很好。
    • 我想我可以使用一个单例对象来缓存每个工人的广播和反序列化值。
    • 我们对大对象所做的不是使用broadcast,而是将对象写入分布式文件系统,然后读取并缓存在执行程序中。缓存使用SoftReference,因此如果内存压力允许,该对象也将在那里用于下一个作业。
    • 你能解释一下如何在不使用某种对象单例的情况下将对象缓存在执行程序中(自 spark 1.2.0 以来这不再工作)?
    • 这是一个单例对象内的简单 HashMap。为什么从 Spark 1.2.0 开始就不行了?
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2016-01-30
    • 1970-01-01
    • 2016-08-18
    • 1970-01-01
    • 1970-01-01
    • 2018-03-11
    相关资源
    最近更新 更多