【发布时间】:2021-12-03 20:07:56
【问题描述】:
spark.memory.offheap.enabled=true 时,Spark 可以利用堆外内存进行洗牌和缓存 (StorageLevel.OFF_HEAP)。堆外内存可以用来存储广播变量吗?怎么样?
【问题讨论】:
标签: apache-spark
spark.memory.offheap.enabled=true 时,Spark 可以利用堆外内存进行洗牌和缓存 (StorageLevel.OFF_HEAP)。堆外内存可以用来存储广播变量吗?怎么样?
【问题讨论】:
标签: apache-spark
简而言之,不,您不能将StorageLevel.OFF_HEAP 用于广播变量。
要了解原因,让我们看看source code 中的SparkContext.broadcast(...) 方法。
/**
* Broadcast a read-only variable to the cluster, returning a
* [[org.apache.spark.broadcast.Broadcast]] object for reading it in distributed functions ...
*/
def broadcast[T: ClassTag](value: T): Broadcast[T] = {
:
val bc = env.broadcastManager.newBroadcast[T](value, isLocal)
:
bc
}
在上面的代码中,broadcastManager.newBroadcast(...) 是创建 Broadcast 对象的原因,它是该方法的返回类型。
现在,让我们深入挖掘并检查newBroadcast()。
def newBroadcast(value_ : T, isLocal: Boolean): Broadcast[T] = {
broadcastFactory.newBroadcast[T](value_, isLocal, nextBroadcastId.getAndIncrement())
}
在上面的代码中,broadcastManager 有一个名为 broadcastFactory 的组件,并使用抽象工厂设计模式将广播变量的创建委托给其相关工厂。
另请注意,BroadcastManager 会跟踪每个 broadcast 变量的唯一 id,该变量会随着每个新的广播变量而递增。
目前spark中可以初始化的BroadcastFactory只有一种,那就是TorrentBroadcastFactory。这可以在BroadcastManager 的initialization code 中看到。
// Called by SparkContext or Executor before using Broadcast
private def initialize() {
:
broadcastFactory = new TorrentBroadcastFactory
:
}
引用TorrentBroadcastFactory的source code
使用类似 BitTorrent 协议的广播实现 将广播数据分布式传输给执行者
这个特定的工厂使用TorrentBroadcast。这个类的描述信息量很大。
驱动将序列化的对象分成小块存储 驱动程序的 BlockManager 中的那些块。
在每个执行器上,执行器首先尝试从其 BlockManager 中获取对象。 如果它不存在,则它使用远程获取来获取小的 如果可用,来自驱动程序和/或其他执行程序的块。一旦它 获取块,它把块放在它自己的 BlockManager 中,准备好 要从中获取的其他执行者。这样可以防止驾驶员 发送多份广播数据的瓶颈 (每个执行者一个)。
阅读TorrentBroadcast类的writeBlock函数,我们可以看到这个广播硬编码的StorageLevel.MEMORY_AND_DISK_SER选项。
/**
* Divide the object into multiple blocks and put those blocks in the block manager.
*
* @param value the object to divide
* @return number of blocks this broadcast variable is divided into
*/
private def writeBlocks(value: T): Int = {
import StorageLevel._
:
:
if (!blockManager.putBytes(pieceId, bytes, MEMORY_AND_DISK_SER, tellMaster = true)) {
throw new SparkException(s"Failed to store $pieceId of $broadcastId " + s"in local BlockManager")
}
:
:
因此,由于此代码使用硬编码值StorageLevel.MEMORY_AND_DISK_SER,我们不能将StorageLevel.OFF_HEAP 用于广播变量。
【讨论】:
MEMORY_AND_DISK,而集群的其余部分——MEMORY_AND_DISK_SER,对吗?
MEMORY_AND_DISK,执行程序使用MEMORY_AND_DISK_SER。我已经更新了答案中的文本,以匹配我在代码 sn-ps 中显示的内容以保持连贯性。