【问题标题】:Spark - using off-heap memorySpark - 使用堆外内存
【发布时间】:2021-12-03 20:07:56
【问题描述】:

spark.memory.offheap.enabled=true 时,Spark 可以利用堆外内存进行洗牌和缓存 (StorageLevel.OFF_HEAP)。堆外内存可以用来存储广播变量吗?怎么样?

【问题讨论】:

    标签: apache-spark


    【解决方案1】:

    简而言之,不,您不能将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 中显示的内容以保持连贯性。
    猜你喜欢
    • 2018-09-27
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2017-10-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多