Spark大数据处理:技术、应用与性能优化【2.0】
4.5 容错机制
在众多特性中,最难实现的是容错性。⼀般来说,分布式数据集的容错性有两种⽅式:数据检查点和记录数据的更新。⾯向⼤规模数据分析,数据检查点操作成本很⾼,需要通过数据中⼼的⽹络连接在机器之间复制庞⼤的数据集,⽽⽹络带宽往往⽐内存带宽低得多,同时还需要消耗更多的存储资源。因此,Spark选择记录更新的⽅式。但是,如果更新粒度太细太多,那么记录更新成本也不低。因此,RDD只⽀持粗粒度转换,即在⼤量记录上执⾏的单个操作。将创建RDD的⼀系列Lineage(即⾎统)记录下来,以便恢复丢失的分区。Lineage本质上很类似于数据库中的重做⽇志(Redo Log),只不过这个重做⽇志粒度很⼤,是对全局数据做同样的重做进⽽恢复数据。
4.5.1 Lineage机制
最后,为了说明模型的容错性,图4-16给出了3个算⼦的⾎统(lineage)关系图。在lines RDD上执⾏filter操作,得到errors,然后filter、map后得到新的RDD(filter、map和collect都是Spark中对RDD的函数操作)。Spark调度器以流⽔线的⽅式执⾏后三个转换,向拥有errors分区缓存的节点发送⼀组任务。此外,如果某个errors分区丢失,则Spark只在相应的lines分区上执⾏filter操作来重建该errors分区。

图4-16 RDD Lineage
1.Lineage简介
相⽐其他系统的细颗粒度的内存数据更新级别的备份或者LOG机制,RDD的Lineage记录的是粗颗粒度的特定数据Transformation操作(如filter、map、join等)⾏为。当这个RDD的部分分区数据丢失时,它可以通过Lineage获取⾜够的信息来重新运算和恢复丢失的数据分区。因为这种粗颗粒的数据模型,限制了Spark的运⽤场合,所以Spark并不适⽤于所有⾼性能要求的场景,但同时相⽐细颗粒度的数据模型,也带来了性能的提升。
2.两种依赖
RDD在Lineage依赖⽅⾯分为两种:Narrow Dependencies与Shuffle Dependencies,⽤来解决数据容错的⾼效性。Narrow Dependencies是指⽗RDD的每⼀个分区最多被⼀个⼦RDD的分区所⽤,表现为⼀个⽗RDD的分区对应于⼀个⼦RDD的分区或多个⽗RDD的分区对应于⼀个⼦RDD的分区,也就是说⼀个⽗RDD的⼀个分区不可能对应⼀个⼦RDD的多个分区。Shuffle Dependencies是指⼦RDD的分区依赖于⽗RDD的多个分区或所有分区,即存在⼀个⽗RDD的⼀个分区对应⼀个⼦RDD的多个分区。
本质理解:根据⽗RDD分区是对应1个还是多个⼦RDD分区来区分Narrow Dependency(⽗分区对应⼀个⼦分区)和Shuffle Dependency (⽗分区对应多个⼦分区)。如果对应多个,则当容错重算分区时,因为⽗分区数据只有⼀部分是需要重算⼦分区的,其余数据重算就造成了冗余计算。
·Narrow Dependency:1个⽗RDD分区对应1个⼦RDD分区,这其中⼜分两种情况:1个⼦RDD分区对应1个⽗RDD分区(如map、filter等算⼦),1个⼦RDD分区对应N个⽗RDD分区(如co-paritioned(协同划分)过的Join)。
·Shuffle Dependency:1个⽗RDD分区对应多个⼦RDD分区,这其中⼜分两种情况:1个⽗RDD对应所有⼦RDD分区(未经协同划分的Join)或者1个⽗RDD对应⾮全部的多个RDD分区(如groupByKey)。
对于Shuffle Dependencies,Stage计算的输⼊和输出在不同的节点上,对于输⼊节点完好,⽽输出节点死机的情况,通过重新计算恢复数据这种情况下,这种⽅法容错是有效的,否则⽆效,因为⽆法重试,需要向上追溯其祖先看是否可以重试(这就是lineage,⾎统的意思),Narrow Dependencies对于数据的重算开销要远⼩于Wide Dependencies的数据重算开销。
Narrow Dependency和Shuffle Dependency的概念主要⽤在两个地⽅:⼀个是容错中相当于Redo⽇志的功能;另⼀个是在调度中构建DAG作为不同Stage的划分点。
3.容错原理
在容错机制中,如果⼀个节点死机了,⽽且运算Narrow Dependency,则只要把丢失的⽗RDD分区重算即可,不依赖于其他节点。⽽Shuffle Dependency需要⽗RDD的所有分区都存在,重算就很昂贵了。可以这样理解开销的经济与否:在Narrow Dependency中,在⼦RDD的分区丢失、重算⽗RDD分区时,⽗RDD相应分区的所有数据都是⼦RDD分区的数据,并不存在冗余计算。在Shuffle Dependency情况下,丢失⼀个⼦RDD分区重算的每个⽗RDD的每个分区的所有数据并不是都给丢失的⼦RDD分区⽤的,会有⼀部分数据相当于对应的是未丢失的⼦RDD分区中需要的数据,这样就会产⽣冗余计算开销,这也是Shuffle Dependency开销更⼤的原因。因此如果使⽤Checkpoint算⼦来做检查点,不仅要考虑Lineage是否⾜够⻓,也要考虑是否有宽依赖,对Shuffle Dependency加Checkpoint是最物有所值的。下⾯结合图4-17进⾏分析。

以图4-17上端的图为例,如果RDD_1中的Partition3出错丢失,则Spark会回溯到Partition3的⽗分区RDD_0的Partition3,对RDD_0的Partition3重算算⼦,得到RDD_1的Partition3。其他分区丢失也是同理重算进⾏容错恢复。
以图4-18下端的图为例,其中RDD_1中的Partition3丢失出错,由于其⽗分区是RDD_0的所有分区,所以需要回溯到RDD_0,重算RDD_0的所有分区,然后将RDD_1的Partition3需要的数据聚集合并为RDD_1的Partition3。在这个过程中,由于RDD_0中不是RDD_1中Partition3需要的数据也全部进⾏了重算,所以产⽣了⼤量冗余数据重算的开销。
通过代码介绍容错的具体调⽤。
下⾯通过CacheManager类的getsOrCompute⽅法作⽤⼊⼝,进⼀步分析容错机制。
def getOrCompute[T](
rdd: RDD[T],
partition: Partition,
context: TaskContext,
storageLevel: StorageLevel): Iterator[T] = {
……
case None =>
val storedValues = acquireLockForPartition[T](key)
if (storedValues.isDefined) {
return new InterruptibleIterator[T](context, storedValues.get)
}
try {
logInfo(s"Partition $key not found, computing it")
/*如果所需的分区丢失了,在BlockManager⽆法找到,这个分区就会重新计算(注意数据⾸次加载也相当于
⽆法找到,需要重新计算)*/
val computedValues = rdd.computeOrReadCheckpoint(partition, context)

下⾯通过RDD的computeOrReadCheckpoint⽅法的代码进⼀步分析容错机制。
private[spark] def computeOrReadCheckpoint(split: Partition, context:
TaskContext): Iterator[T] =
{
/*这⾥相当于对这个分区回溯到⽗节点或者祖先节点,然后⼀路计算回来得到这个分区,相当于只需要计算这
个分区的依赖,因为是获取这个分区,⽽不是计算所有分区*/
( ) ( )if (isCheckpointed) firstParent[T].iterator(split, context) else compute
(split, context)
}
可以通过图4-19来理解重做Lineage的过程,虚线⽅框表⽰逻辑分区,相当于计算完成后就不存在了,已经转化为RDD_2中分区的数据。实线⽅框表⽰分区的重新计算过程就是由于RDD_2的分区丢失了,程序⽤到Partition0分区,找不到,就反向回溯Lineage到RDD_0的分区Partition0和Partition1,然后对其进⾏重新计算,计算结果为RDD_1的Partition0,RDD_1再重新计算Partition0为RDD_2中的Partition0。这时就不需要其他分区参与计算了。

4.5.2 Checkpoint机制
通过上述分析可以看出在以下两种情况下,RDD需要加检查点。
1)DAG中的Lineage过⻓,如果重算,则开销太⼤(如在PageRank中)。
2)在Shuffle Dependency上做Checkpoint(检查点)获得的收益更⼤。
由于RDD是只读的,所以Spark的RDD计算中⼀致性不是主要关⼼的内容,内存相对容易管理,这也是设计者很有远⻅的地⽅,这样减少了框架的复杂性,提升了性能和可扩展性,为以后上层框架的丰富奠定了强有⼒的基础。
在RDD计算中,通过检查点机制进⾏容错,传统做检查点有两种⽅式:通过冗余数据和⽇志记录更新操作。在RDD中的doCheckPoint⽅法相当于通过冗余数据来缓存数据,⽽之前介绍的⾎统就是通过相当粗粒度的记录更新操作来实现容错的。
在Spark中,通过RDD中的checkpoint()⽅法来做检查点。
def checkpoint():Unit
可以通过SparkContext.setCheckPointDir()设置检查点数据的存储路径,进⽽将数据存储备份,然后Spark删除所有已经做检查点的RDD的祖先RDD依赖。这个操作需要在所有需要对这个RDD所做的操作完成之后再做,因为数据会写⼊持久化存储造成I/O开销。官⽅建议,做检查点的RDD最好是在内存中已经缓存的RDD,否则保存这个RDD在持久化的⽂件中需要重新计算,产⽣I/O开销。
下⾯通过源码来了解检查点的机制。
检查点(本质是通过将RDD写⼊Disk做检查点)是为了通过lineage做容错的辅助,lineage过⻓会造成容错成本过⾼,这样就不如在中间阶段做检查点容错,如果之后有节点出现问题⽽丢失分区,从做检查点的RDD开始重做Lineage,就会减少开销。
在RDD中通过doCheckpoint()⽅法作为检查点的⼊⼝⽅法。
private[spark] def doCheckpoint() { ……
checkpointData.get.doCheckpoint() } else
{ dependencies.foreach(_.rdd.doCheckpoint()) } …… }
在RDDCheckpointData中,通过doCheckpoint()⽅法做检查点。
def doCheckpoint() { ……
RDD通过同步⽅式做检查点,具体使⽤Synchronized保证⽅法的同步和线程安全。代码实现如下
/*path检查点RDD的输出⽂件路径*/
CheckpointData.synchronized { ……
val path = new Path(rdd.context.checkpointDir.get, "rdd-" + rdd.id)
val fs = path.getFileSystem(new Configuration())
if (!fs.mkdirs(path)) {
throw new SparkException("Failed to create checkpoint path " + path) }
/*在SparkContext提交作业,将检查点RDD写⼊之前设置的路径中*/
rdd.context.runJob(rdd, CheckpointRDD.writeToFile(path.toString) _)
val newRDD = new CheckpointRDD[T](rdd.context, path.toString) ……
}
/*在CheckPointRDD中调⽤ writeToFile⽅法将RDD写⼊HDFS*/
def writeToFile[T](path: String, blockSize: Int = -1)(ctx:
TaskContext, iterator: Iterator[T]) {
val env = SparkEnv.get
val outputDir = new Path(path)
/*本质相当于在Hadoop 的分布式⽂件系统将RDD数据写进HDFS*/
val fs = outputDir.getFileSystem(env.hadoop.newConfiguration()) val
finalOutputName = splitIdToFile(ctx.splitId)
val finalOutputPath = new Path(outputDir, finalOutputName)
val tempOutputPath = new Path(outputDir, "." + finalOutputName + "-attempt-" +
ctx.attemptId) ……
/*根据数据量不同设置,不同的缓冲区⼤⼩*/
val bufferSize = System.getProperty("spark.buffer.size", "65536").toInt
val fileOutputStream = if (blockSize < 0) {
fs.create(tempOutputPath, false, bufferSize)
} else { // This is mainly for testing purpose
fs.create(tempOutputPath, false, bufferSize,
fs.getDefaultReplication, blockSize) }
/*创建序列化器*/
val serializer = env.serializer.newInstance()
val serializeStream = serializer.serializeStream(fileOutputStream)
/*此处为写⼊操作,关键是在iterator上相当于将iteraor迭代器的对象序列化写到HDFS中*/
serializeStream.writeAll(iterator) serializeStream.close() …… }
SerializationStream写⼊操作
trait SerializationStream {
def writeObject[T](t:T): SerializationStream
def flush():Unit
def close():Unit
def writeAll[T](iter: Iterator[T]): SerializationStream = {
while (iter.hasNext) { writeObject(iter.next()) } this } }
/*如果配置了kyro序列化器进⾏写⼊,则调⽤下⾯的writeObject⽅法将数据序列化后写⼊HDFS*/
private[spark] class KryoSerializationStream(kryo: Kryo, outStream:
OutputStream) extends SerializationStream
{ val output = new KryoOutput(outStream)
……
def writeObject[T](t: T): SerializationStream = {
kryo.writeClassAndObject(output, t)
this }
def flush() { output.flush() }
def close() { output.close() }
……
}
更多推荐
所有评论(0)