Advertisement

磁盘块管理器

阅读量:

文章结构概览

  • 简介
    • 建立
    • 核心成员
    • 主要作用
      • 在本地系统中生成文件夹
      • 确定与BlockId相对应的文件存储路径
      • 检查BlockId所对应的文件是否存在于系统中
      • 生成临时存储Block数据的文件

简介

DiskBlockManager的核心功能是建立逻辑blocks块与物理磁盘位置之间的对应关系。每个逻辑block块均通过其唯一的BlockId标识,进而被映射至磁盘中的特定文件。这些block块所对应的文件按照哈希算法进行存储,并被放置在由spark.local.dir(或SPARK_LOCAL_DIRS)配置参数所指定的目录列表中。

创建

在DiskBlockManager的初始化过程中,系统将生成对应的BlockManager实例:

复制代码
    val diskBlockManager = {
      // Only perform cleanup if an external service is not serving our shuffle files.
      val deleteFilesOnStop =
    !externalShuffleServiceEnabled || executorId == SparkContext.DRIVER_IDENTIFIER
      new DiskBlockManager(conf, deleteFilesOnStop)
    }
    
    
      
      
      
      
      
      
    

主要成员

  • localDirs: Array[File]:依据spark.local.dir配置参数所指定的路径,生成对应的目录集合。随后在这些主目录中建立用于存储临时数据的子文件夹(此举旨在防止顶层目录的inodes数量过于庞大),最终将文件按照其名称的哈希值分配至不同的子目录中。
    • subDirsPerLocalDir:表示每个由localDirs所代表的主目录下可创建的子目录数量,该数值由配置项 spark.diskStore.subDirectories决定,默认设定为64个。
    • subDirs: Array[Array[File]]:表示每个主目录下的子目录结构,其数量由subDirsPerLocalDir参数确定。该数组本身具有不可变性,但其中每一个元素subDirs(i)是可修改的,并且每个子目录均受到自身锁机制的保护。
    • shutdownHook:即进程终止时触发的关闭钩子函数,用于递归删除所有归属于当前Application实例的、位于localDirs路径下的文件。

主要功能概述

DiskBlockManager负责维护每个BlockId与磁盘文件之间的对应关系。该组件具备以下主要功能:

  • 建立本地存储目录;
  • 确定特定BlockId所对应的文件存储路径;
  • 检查某BlockId对应的文件是否存在于磁盘中;
  • 生成临时Block文件。

创建本地目录

在本地节点上,负责配置创建的特定列表磁盘目录被用于将Block数据保存至指定文件中。当启用外部shuffle服务时,JVM退出后该目录将不会被清除。

复制代码
    private def createLocalDirs(conf: SparkConf): Array[File] = {
      // 获取配置的本地目录列表
      Utils.getConfiguredLocalDirs(conf).flatMap { rootDir =>
    try {
    	  // 创建目录,目录名为rootDir + "blockmgr-" + UUID.randomUUID.toString
      val localDir = Utils.createDirectory(rootDir, "blockmgr")
      logInfo(s"Created local directory at $localDir")
      Some(localDir)
    } catch {
      case e: IOException =>
        logError(s"Failed to create local dir in $rootDir. Ignoring this directory.", e)
        None
    }
      }
    }
    
    
      
      
      
      
      
      
      
      
      
      
      
      
      
      
      
    

获取BlockId对应的文件路径

当Block数据需要持久化存储时,首要操作是调用getFile方法以分配一个专属的文件路径(同样通过该方法可定位到对应blockId的文件)。其具体处理流程如下:

首先,依据文件名生成对应的哈希值;

随后,将该哈希值与本地一级目录的数量进行取模运算,得到的结果记为dirId;

接着,再次利用哈希值与一级目录总数进行除法运算,所得商数再与二级目录的数量取模,结果记为subDirId;

若dirId/subDirId对应的目录已存在,则读取该目录下的文件;若不存在,则创建新的dirId/subDirId目录。

复制代码
      def getFile(blockId: BlockId): File = getFile(blockId.name)

    
      // 通过将文件名hash到本地子目录的方式来查找文件
      // 该方法应该与"org.apache.spark.network.shuffle.ExternalShuffleBlockResolver#getFile()"方法保持同步
      def getFile(filename: String): File = {
    // 得到该文件名hash到的本地目录,以及该目录下的子目录
    val hash = Utils.nonNegativeHash(filename)
    val dirId = hash % localDirs.length
    val subDirId = (hash / localDirs.length) % subDirsPerLocalDir
      
    // 如果目录不存在则创建该子目录
    val subDir = subDirs(dirId).synchronized {
      val old = subDirs(dirId)(subDirId)
      if (old != null) {
        old
      } else {
        val newDir = new File(localDirs(dirId), "%02x".format(subDirId))
        if (!newDir.exists() && !newDir.mkdir()) {
          throw new IOException(s"Failed to create local dir in $newDir.")
        }
        subDirs(dirId)(subDirId) = newDir
        newDir
      }
    }
      
    new File(subDir, filename)
      }
    
    
        
        
        
        
        
        
        
        
        
        
        
        
        
        
        
        
        
        
        
        
        
        
        
        
        
        
        

查询BlockId对应的文件是否存在

判断特定BlockId所关联的文件是否存在,实际上是通过核查该BlockId对应的文件路径是否真实存在来完成的。

复制代码
    /** Check if disk block manager has a block. */
    def containsBlock(blockId: BlockId): Boolean = {
      getFile(blockId.name).exists()
    }
    
    
      
      
      
      
    

创建临时Block文件

在ShuffleMapTask执行完毕后,为存储中间结果需要调用createTempShuffleBlock方法生成临时的Block文件,并返回一个TempShuffleBlockId。该标识符的生成方式为在字符串"temp_shuffle_"后拼接一个UUID,同时与对应的文件进行关联。

复制代码
    /** Produces a unique block id and File suitable for storing shuffled intermediate results. */
    def createTempShuffleBlock(): (TempShuffleBlockId, File) = {
      var blockId = new TempShuffleBlockId(UUID.randomUUID())
      while (getFile(blockId).exists()) {
    blockId = new TempShuffleBlockId(UUID.randomUUID())
      }
      (blockId, getFile(blockId))
    }
    
    /** Id associated with temporary shuffle data managed as blocks. Not serializable. */
    private[spark] case class TempShuffleBlockId(id: UUID) extends BlockId {
      override def name: String = "temp_shuffle_" + id
    }
    
    
      
      
      
      
      
      
      
      
      
      
      
      
      
    

全部评论 (0)

还没有任何评论哟~