如何在Spark中从堆中删除/处理广播变量?

sam*_*est 18 memory-management scala apache-spark

要广播一个变量,使得一个变量在一个集群上的每个节点的内存中只发生一次,就能做到:val myVarBroadcasted = sc.broadcast(myVar)然后在RDD转换中检索它,如下所示:

myRdd.map(blar => {
  val myVarRetrieved = myVarBroadcasted.value
  // some code that uses it
}
.someAction
Run Code Online (Sandbox Code Playgroud)

但是现在假设我希望用新的广播变量执行更多操作 - 如果由于旧的广播变量而没有足够的堆空间怎么办?我想要一个像这样的功能

myVarBroadcasted.remove()
Run Code Online (Sandbox Code Playgroud)

现在我似乎找不到这样做的方法.

另外,一个非常相关的问题:广播变量在哪里?它们会进入总内存的缓存分数,还是只进入堆分数?

Gia*_*gna 26

如果要从必须使用的执行程序和驱动程序中删除广播变量destroy,请unpersist仅使用从执行程序中删除它:

myVarBroadcasted.destroy()
Run Code Online (Sandbox Code Playgroud)

这种方法是封锁的.我喜欢面食!


Shy*_*nki 11

您正在寻找Spark 1.0.0 提供的unpersist

myVarBroadcasted.unpersist(blocking = true)
Run Code Online (Sandbox Code Playgroud)

广播变量存储为反序列化Java对象或序列化ByteBuffers的ArrayBuffers.(存储方面,它们被视为类似于RDD - 需要确认)

unpersist方法将它们从内存以及每个执行程序节点上的磁盘中删除.但它保留在驱动程序节点上,因此可以重新广播.