磁盘块管理器
文章结构概览
- 简介
- 建立
- 核心成员
- 主要作用
-
- 在本地系统中生成文件夹
- 确定与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
}
