阅读背景:

第223讲:Spark Shuffle Pluggable框架ShuffleReader解析

来源:互联网 
在reduce任务中,读取mappers中的聚合数据。
/**   * Interface to get local block data. Throws an exception if the block cannot be found or   * cannot be read successfully.   */  override def getBlockData(blockId: BlockId): ManagedBuffer = {    if (blockId.isShuffle) {      shuffleManager.shuffleBlockResolver.getBlockData(blockId.asInstanceOf[ShuffleBlockId])    } else {      val blockBytesOpt = doGetLocal(blockId, asBlockResult = false)        .asInstanceOf[Option[ByteBuffer]]      if (blockBytesOpt.isDefined) {        val buffer = blockBytesOpt.get        new NioManagedBuffer(buffer)      } else {        throw new BlockNotFoundException(blockId.toString)      }    }  }/**   * Interface to g



你的当前访问异常,请进行认证后继续阅读剩余内容。

分享到: