zoukankan      html  css  js  c++  java
  • spark中的广播变量broadcast

    Spark中的Broadcast处理

    首先先来看一看broadcast的使用代码:

    val values = List[Int](1,2,3)

    val broadcastValues = sparkContext.broadcast(values)

    rdd.mapPartitions(iter => {

      broadcastValues.getValue.foreach(println)

    })

     

    在上面的代码中,首先生成了一个集合变量,把这个变量通过sparkContext的broadcast函数进行广播,

    最后在rdd的每个partition的迭代时,使用这个广播变量.

     

    接下来看看广播变量的生成与数据的读取实现部分:

    def broadcast[T: ClassTag](value: T): Broadcast[T] = {
      assertNotStopped()
      if (classOf[RDD[_]].isAssignableFrom(classTag[T].runtimeClass)) {

    这里要注意,使用broadcast时,不能直接对RDD进行broadcast的操作.
        // This is a warning instead of an exception in order to avoid breaking

    //       user programs that
        // might have created RDD broadcast variables but not used them:
        logWarning("Can not directly broadcast RDDs; instead, call collect() and "
          "broadcast the result (see SPARK-5063)")
      }

     

    通过broadcastManager中的newBroadcast函数来进行广播.
      val bc = env.broadcastManager.newBroadcast[T](valueisLocal)
      val callSite = getCallSite
      logInfo("Created broadcast " + bc.id + " from " + callSite.shortForm)
      cleaner.foreach(_.registerBroadcastForCleanup(bc))
      bc
    }

     

    在BroadcastManager中生成广播变量的函数,这个函数直接使用的broadcastFactory的相应函数.

    broadcastFactory的实例通过配置spark.broadcast.factory,

         默认是TorrentBroadcastFactory.

    def newBroadcast[T: ClassTag](value_ : TisLocal: Boolean): Broadcast[T] = {
      broadcastFactory.newBroadcast[T](value_isLocal

           nextBroadcastId.getAndIncrement())
    }

     

    在TorrentBroadcastFactory中生成广播变量的函数:

    在这里面,直接生成了一个TorrentBroadcast的实例.

    override def newBroadcast[T: ClassTag](value_ : TisLocal: Boolean, id: Long)

    : Broadcast[T] = {
      new TorrentBroadcast[T](value_id)
    }

     

    TorrentBroadcast实例生成时的处理流程:

    这里基本的代码部分是直接写入这个要广播的变量,返回的值是这个变量所占用的block的个数.

    Broadcast的block的大小通过spark.broadcast.blockSize配置.默认是4MB,

    Broadcast的压缩是否通过spark.broadcast.compress配置,默认是true表示启用,默认情况下使用snappy的压缩.

     

    private val broadcastId BroadcastBlockId(id)
    /** Total number of blocks this broadcast variable contains. */
    private val numBlocksInt = writeBlocks(obj)

     

    接下来生成一个lazy的属性,这个属性仅仅有在详细的使用时,才会运行,在实例生成时不运行(上面的演示样例中的getValue.foreach时运行).

    @transient private lazy val _value= readBroadcastBlock()

    override protected def getValue() = {
      _value
    }

     

    看看实例生成时的writeBlocks的函数:

    private def writeBlocks(value: T): Int = {

    这里先把这个广播变量保存一份到当前的task的storage中,这样做是保证在读取时,假设要使用这个广播变量的task就是本地的task时,直接从blockManager中本地读取.
      SparkEnv.get.blockManager.putSingle(broadcastIdvalue

    StorageLevel.MEMORY_AND_DISK,
        tellMaster = false)

     

    这里依据block的设置大小,对value进行序列化/压缩分块,每个块的大小为blocksize的大小,
      val blocks =
        TorrentBroadcast.blockifyObject(valueblockSizeSparkEnv.get.serializer

        compressionCodec)

     

    这里把序列化并压缩分块后的blocks进行迭代,存储到blockManager中,
      blocks.zipWithIndex.foreach { case (blocki) =>
        SparkEnv.get.blockManager.putBytes(
          BroadcastBlockId(id"piece" + i),
          block,
          StorageLevel.MEMORY_AND_DISK_SER,
          tellMaster = true)
      }

    这个函数的返回值是一个int类型的值,这个值就是序列化压缩存储后block的个数.
      blocks.length
    }

     

    在我们的演示样例中,使用getValue时,会运行实例初始化时定义的lazy的函数readBroadcastBlock:

    private def readBroadcastBlock(): = Utils.tryOrIOException {
      TorrentBroadcast.synchronized {
        setConf(SparkEnv.get.conf)

    这里先从local端的blockmanager中直接读取storage中相应此广播变量的内容,假设能读取到,表示这个广播变量已经读取过来或者说这个task就是广播的本地executor.
        SparkEnv.get.blockManager.getLocal(broadcastId).map(_.data.next()) match {
          case Some(x) =>
            x.asInstanceOf[T]

    以下这部分运行时,表示这个广播变量在当前的executor中是第一次读取,通过readBlocks函数去读取这个广播变量的全部的blocks,反序列化后,直接把这个广播变量存储到本地的blockManager中,下次读取时,就能够直接从本地进行读取.
          case None =>
            logInfo("Started reading broadcast variable " + id)
            val startTimeMs = System.currentTimeMillis()
            val blocks = readBlocks()
            logInfo("Reading broadcast variable " + id + " took" 

                  Utils.getUsedTimeMs(startTimeMs))

            val obj = TorrentBroadcast.unBlockifyObject[T](
              blocksSparkEnv.get.serializercompressionCodec)
            // Store the merged copy in BlockManager so other tasks on this executor don't
            // need to re-fetch it.
            SparkEnv.get.blockManager.putSingle(
              broadcastIdobjStorageLevel.MEMORY_AND_DISKtellMaster = false)
            obj
        }
      }
    }

     

    最后再看看readBlocks函数的处理流程:

    private def readBlocks(): Array[ByteBuffer] = {

    这里定义的变量用于存储读取到的block的信息,numBlocks是广播变量序列化后所占用的block的个数.
      val blocks = new Array[ByteBuffer](numBlocks)
      val bm = SparkEnv.get.blockManager

    这里開始迭代读取每个block的内容,这里的读取是先从local中进行读取,假设local中没有读取到数据时,通过blockManager读取远端的数据,通过读取这个block相应的location从这个location去读取这个block的内容,并存储到本地的blockManager中.最后,这个函数返回读取到的blocks的集合.
      for (pid <- Random.shuffle(Seq.range(0numBlocks))) {
        val pieceId = BroadcastBlockId(id"piece" + pid)
        logDebug(s"Reading piece $pieceId of $broadcastId")

        def getLocal: Option[ByteBuffer] = bm.getLocalBytes(pieceId)
        def getRemote: Option[ByteBuffer] = bm.getRemoteBytes(pieceId).map { block =>
          SparkEnv.get.blockManager.putBytes(
            pieceId,
            block,
            StorageLevel.MEMORY_AND_DISK_SER,
            tellMaster = true)
          block
        }
        val block: ByteBuffer = getLocal.orElse(getRemote).getOrElse(
          throw new SparkException(s"Failed to get $pieceId of $broadcastId"))
        blocks(pid) = block
      }
      blocks
    }

  • 相关阅读:
    【转】Redis概念原理、redis面试
    mysql登录后显示用户名与当前数据库名
    (5.3.10)数据库迁移——sql server降级操作
    navicat下载安装破解
    sql server2016+windows server2016使用日志传送做主从,主库无法备份事务日志
    (4.41)sql server如何把xml转换成表格数据?
    .NET Core SDK在Windows系统安装后出现Failed to load the hostfxr.dll等问题的解决方法
    .NET Core实战项目之CMS 第七章 设计篇-用户权限极简设计全过程
    .NET Core实战项目之CMS 第六章 入门篇-Vue的快速入门及其使用
    [译]聊聊C#中的泛型的使用(新手勿入)
  • 原文地址:https://www.cnblogs.com/claireyuancy/p/7057675.html
Copyright © 2011-2022 走看看