分享

Spark Core源码分析: RDD基础

52Pig 2014-10-17 23:07:28 发表于 代码分析 [显示全部楼层] 回帖奖励 阅读模式 关闭右栏 1 13749

阅读导读:
1.getPartitions和compute进行了哪些操作?
2.hadoop如何进行序列化?
3.checkpoint的执行逻辑?





RDD

RDD初始参数:上下文和一组依赖


  1. abstract class RDD[T: ClassTag](
  2.     @transient private var sc: SparkContext,
  3.     @transient private var deps: Seq[Dependency[_]]
  4.   ) extends Serializable
复制代码

以下需要仔细理清:

A list of Partitions

Function to compute split (sub RDD impl)

A list of Dependencies

Partitioner for K-V RDDs (Optional)

Preferred locations to compute each spliton (Optional)

Dependency

Dependency代表了RDD之间的依赖关系,即血缘


RDD中的使用RDD给子类提供了getDependencies方法来制定如何依赖父类RDD


  1. protected def getDependencies: Seq[Dependency[_]] = deps
复制代码

事实上,在获取first parent的时候,子类经常会使用下面这个方法
  1. protected[spark] def firstParent[U: ClassTag] = {
  2.   dependencies.head.rdd.asInstanceOf[RDD[U]]
  3. }
复制代码
可以看到,Seq里的第一个dependency应该是直接的parent,从而从第一个dependency类里获得了rdd,这个rdd就是父RDD。
一般的RDD子类都会这么实现compute和getPartition方法,以SchemaRDD举例:
  1. override def compute(split: Partition, context: TaskContext): Iterator[Row] =
  2.     firstParent[Row].compute(split, context).map(_.copy())
  3. override def getPartitions: Array[Partition] = firstParent[Row].partitions
复制代码
compute()方法调用了第一个父类的compute,把结果RDD copy返回,getPartitions返回的就是第一个父类的partitions
下面看一下Dependency类及其子类的实现。
  1. abstract class Dependency[T](val rdd: RDD[T]) extends Serializable
复制代码
Dependency里传入的rdd,就是父RDD本身。继承结构如下:
  1. def getParents(partitionId: Int): Seq[Int]
复制代码
OneToOneDependency表示父RDD和子RDD的分区依赖是一对一的。
RangeDependency表示在一个range范围内,依赖关系是一对一的,所以初始化的时候会有一个范围,范围外的partitionId,传进去之后返回的是Nil。
下面介绍宽依赖。
  1. class ShuffleDependency[K, V](
  2.     @transient rdd: RDD[_ <: Product2[K, V]],
  3.     val partitioner: Partitioner,
  4.     val serializer: Serializer = null)
  5.   extends Dependency(rdd.asInstanceOf[RDD[Product2[K, V]]]) {
  6.   // 上下文增量定义的Id
  7.   val shuffleId: Int = rdd.context.newShuffleId()
  8.   // ContextCleaner的作用和实现在SparkContext章节叙述
  9.   rdd.sparkContext.cleaner.foreach(_.registerShuffleForCleanup(this))
  10. }
复制代码
宽依赖针对的RDD是KV形式的,需要一个partitioner指定分区方式(下一节介绍),需要一个序列化工具类,序列化工具目前的实现如下:

宽依赖和窄依赖对失败恢复时候的recompute有不同程度的影响,宽依赖可能是要全部计算的。
Partition

Partition具体表示RDD每个数据分区。

Partition提供trait类,内含一个index和hashCode()方法,具体子类实现与RDD子类有关,种类如下:



在分析每个RDD子类的时候再涉及。
PartitionerPartitioner决定KV形式的RDD如何根据key进行partition



  1. abstract class Partitioner extends Serializable {
  2.   def numPartitions: Int // 总分区数
  3.   def getPartition(key: Any): Int
  4. }
复制代码
在ShuffleDependency里对应一个Partitioner,来完成宽依赖下,子RDD如何获取父RDD。
默认PartitionerPartitioner的伴生对象提供defaultPartitioner方法,逻辑为:
传入的RDD(至少两个)中,遍历(顺序是partition数目从大到小)RDD,如果已经有Partitioner了,就使用。如果RDD们都没有Partitioner,则使用默认的HashPartitioner。而HashPartitioner的初始化partition数目,取决于是否设置了spark.default.parallelism,如果没有的话就取RDD中partition数目最大的值。


如果上面这段文字看起来费解,代码如下:
  1.   def defaultPartitioner(rdd: RDD[_], others: RDD[_]*): Partitioner = {
  2.     val bySize = (Seq(rdd) ++ others).sortBy(_.partitions.size).reverse
  3.     for (r <- bySize if r.partitioner.isDefined) {
  4.       return r.partitioner.get
  5.     }
  6.     if (rdd.context.conf.contains("spark.default.parallelism")) {
  7.       new HashPartitioner(rdd.context.defaultParallelism)
  8.     } else {
  9.       new HashPartitioner(bySize.head.partitions.size)
  10.     }
  11.   }
复制代码
HashPartitionerHashPartitioner基于java的Object.hashCode。会有个问题是Java的Array有自己的hashCode,不基于Array里的内容,所以RDD[Array[_]]或RDD[(Array[_], _)]使用HashPartitioner会有问题。

顾名思义,getPartition方法实现如下:
  1.   def getPartition(key: Any): Int = key match {
  2.     case null => 0
  3.     case _ => Utils.nonNegativeMod(key.hashCode, numPartitions)
  4.   }
复制代码
RangePartitionerRangePartitioner处理的KV RDD要求Key是可排序的,即满足Scala的Ordered[K]类型。所以它的构造如下:

  1. class RangePartitioner[K <% Ordered[K]: ClassTag, V](
  2.     partitions: Int,
  3.     @transient rdd: RDD[_ <: Product2[K,V]],
  4.     private val ascending: Boolean = true)
  5.   extends Partitioner {
复制代码
内部会计算一个rangBounds(上界),在getPartition的时候,如果rangBoundssize小于1000,则逐个遍历获得;否则二分查找获得partitionId。
Persist默认cache()过程是将RDD persist在内存里,persist()操作可以为RDD重新指定StorageLevel


  1. class StorageLevel private(
  2.     private var useDisk_ : Boolean,
  3.     private var useMemory_ : Boolean,
  4.     private var useOffHeap_ : Boolean,
  5.     private var deserialized_ : Boolean,
  6.     private var replication_ : Int = 1)
复制代码
  1. object StorageLevel {
  2.   val NONE = new StorageLevel(false, false, false, false)
  3.   val DISK_ONLY = new StorageLevel(true, false, false, false)
  4.   val DISK_ONLY_2 = new StorageLevel(true, false, false, false, 2)
  5.   val MEMORY_ONLY = new StorageLevel(false, true, false, true)
  6.   val MEMORY_ONLY_2 = new StorageLevel(false, true, false, true, 2)
  7.   val MEMORY_ONLY_SER = new StorageLevel(false, true, false, false)
  8.   val MEMORY_ONLY_SER_2 = new StorageLevel(false, true, false, false, 2)
  9.   val MEMORY_AND_DISK = new StorageLevel(true, true, false, true)
  10.   val MEMORY_AND_DISK_2 = new StorageLevel(true, true, false, true, 2)
  11.   val MEMORY_AND_DISK_SER = new StorageLevel(true, true, false, false)
  12.   val MEMORY_AND_DISK_SER_2 = new StorageLevel(true, true, false, false, 2)
  13.   val OFF_HEAP = new StorageLevel(false, false, true, false) // Tachyon
复制代码
RDD的persist()和unpersist()操作,都是由SparkContext执行的(SparkContext的persistRDD和unpersistRDD方法)。

Persist过程是把该RDD存在上下文的TimeStampedWeakValueHashMap里维护起来。也就是说,其实persist并不是action,并不会触发任何计算。

Unpersist过程如下,会交给SparkEnv里的BlockManager处理。

  1.   private[spark] def unpersistRDD(rddId: Int, blocking: Boolean = true) {
  2.     env.blockManager.master.removeRdd(rddId, blocking)
  3.     persistentRdds.remove(rddId)
  4.     listenerBus.post(SparkListenerUnpersistRDD(rddId))
  5.   }
复制代码
Checkpoint
RDD Actions api里提供了checkpoint()方法,会把本RDD save到SparkContext CheckpointDir
目录下。建议该RDD已经persist在内存中,否则需要recomputation。

如果该RDD没有被checkpoint过,则会生成新的RDDCheckpointData。RDDCheckpointData类与一个RDD关联,记录了checkpoint相关的信息,并且记录checkpointRDD的一个状态,

[ Initialized --> marked for checkpointing--> checkpointing in progress --> checkpointed ]

内部有一个doCheckpoint()方法(会被下面调用)。


执行逻辑真正的checkpoint触发,在RDD私有方法doCheckpoint()里。doCheckpoint()会被DAGScheduler调用,且是在此次job里使用这个RDD完毕之后,此时这个RDD就已经被计算或者物化过了。可以看到,会对RDD的父RDD进行递归。


  1.   private[spark] def doCheckpoint() {
  2.     if (!doCheckpointCalled) {
  3.       doCheckpointCalled = true
  4.       if (checkpointData.isDefined) {
  5.         checkpointData.get.doCheckpoint()
  6.       } else {
  7.         dependencies.foreach(_.rdd.doCheckpoint())
  8.       }
  9.     }
  10.   }
复制代码
RDDCheckpointData的doCheckpoint()方法关键代码如下:
  1. // Create the output path for the checkpoint
  2. val path = new Path(rdd.context.checkpointDir.get, "rdd-" + rdd.id)
  3. val fs = path.getFileSystem(rdd.context.hadoopConfiguration)
  4. if (!fs.mkdirs(path)) {
  5.   throw new SparkException("Failed to create checkpoint path " + path)
  6. }
  7. // Save to file, and reload it as an RDD
  8. val broadcastedConf = rdd.context.broadcast(
  9.   new SerializableWritable(rdd.context.hadoopConfiguration))
  10. // 这次runJob最终调的是dagScheduler的runJob
  11. rdd.context.runJob(rdd,
  12. CheckpointRDD.writeToFile(path.toString, broadcastedConf) _)
  13. // 此时rdd已经记录到磁盘上
  14. val newRDD = new CheckpointRDD[T](rdd.context, path.toString)
  15. if (newRDD.partitions.size != rdd.partitions.size) {
  16.   throw new SparkException("xxx")
  17. }
复制代码

runJob最终调的是dagScheduler的runJob。做完后,生成一个CheckpointRDD。

具体CheckpointRDD相关内容可以参考其他章节。

API子类需要实现的方法


  1. // 计算某个分区
  2. def compute(split: Partition, context: TaskContext): Iterator[T]
  3. protected def getPartitions: Array[Partition]
  4. // 依赖的父RDD,默认就是返回整个dependency序列
  5. protected def getDependencies: Seq[Dependency[_]] = deps
  6. protected def getPreferredLocations(split: Partition): Seq[String] = Nil
复制代码
Transformations
略。
Actions略。
SubRDDs部分RDD子类的实现分析,包括以下几个部分:
1)  子类本身构造参数
2)  子类的特殊私有变量
3)  子类的Partitioner实现
4)  子类的父类函数实现


  1. def compute(split: Partition, context: TaskContext): Iterator[T]
  2. protected def getPartitions: Array[Partition]
  3. protected def getDependencies: Seq[Dependency[_]] = deps
  4. protected def getPreferredLocations(split: Partition): Seq[String] = Nil
复制代码
CheckpointRDD
  1. class CheckpointRDD[T: ClassTag](sc: SparkContext, val checkpointPath: String)
  2.   extends RDD[T](sc, Nil)
复制代码
CheckpointRDDPartition继承自Partition,没有什么增加。
有一个被广播的hadoop conf变量,在compute方法里使用(readFromFile的时候用)
  1. val broadcastedConf = sc.broadcast(
  2. new SerializableWritable(sc.hadoopConfiguration))
复制代码
getPartitions: Array[Partition]方法:
根据checkpointPath去查看Path下有多少个partitionFile,File个数为partition数目。getPartitions方法返回的Array[Partition]内容为New CheckpointRDDPartition(i),i为[0, 1, …, partitionNum]
getPreferredLocations(split:Partition): Seq[String]方法:
文件位置信息,借助hadoop core包,获得block location,把得到的结果按照host打散(flatMap)并过滤掉localhost,返回。
compute(split: Partition, context:TaskContext): Iterator[T]方法:
调用CheckpointRDD.readFromFile(file, broadcastedConf,context)方法,其中file为hadoopfile path,conf为广播过的hadoop conf。

Hadoop文件读写及序列化伴生对象提供writeToFile方法和readFromFile方法,主要用于读写hadoop文件,并且利用env下的serializer进行序列化和反序列化工作。两个方法具体实现如下:

  1. def writeToFile[T](
  2. path: String,
  3. broadcastedConf: Broadcast[SerializableWritable[Configuration]],
  4. blockSize: Int = -1
  5. )(ctx: TaskContext, iterator: Iterator[T]) {
复制代码
创建hadoop文件的时候会若存在会抛异常。把hadoop的outputStream放入serializer的stream里,serializeStream.writeAll(iterator)写入。
writeToFile的调用在RDDCheckpointData类的doCheckpoint方法里,如下:
  1. rdd.context.runJob(rdd,
  2. CheckpointRDD.writeToFile(path.toString, broadcastedConf) _)
复制代码
  1. def readFromFile[T](
  2.   path: Path,
  3.   broadcastedConf: Broadcast[SerializableWritable[Configuration]],
  4.   context: TaskContext
  5. ): Iterator[T] = {
复制代码

打开Hadoop的inutStream,读取的时候使用env下的serializer得到反序列化之后的流。返回的时候,DeserializationStream这个trait提供了asIterator方法,每次next操作可以进行一次readObject。

在返回之前,调用了TaskContext提供的addOnCompleteCallback回调,用于关闭hadoop的inputStream。

NewHadoopRDD
  1. class NewHadoopRDD[K, V](
  2.     sc : SparkContext,
  3.     inputFormatClass: Class[_ <: InputFormat[K, V]],
  4.     keyClass: Class[K],
  5.     valueClass: Class[V],
  6.     @transient conf: Configuration)
  7.   extends RDD[(K, V)](sc, Nil)
  8.   with SparkHadoopMapReduceUtil
复制代码
  1. private[spark] class NewHadoopPartition(
  2.     rddId: Int,
  3.     val index: Int,
  4.     @transient rawSplit: InputSplit with Writable)
  5.   extends Partition {
  6.   val serializableHadoopSplit = new SerializableWritable(rawSplit)
  7.   override def hashCode(): Int = 41 * (41 + rddId) + index
  8. }
复制代码

getPartitions操作:

根据inputFormatClass和conf,通过hadoop InputFormat实现类的getSplits(JobContext)方法得到InputSplits。(ORCFile在此处的优化)

这样获得的split同RDD的partition直接对应。

compute操作:

针对本次split(partition),调用InputFormat的createRecordReader(split)方法,

得到RecordReader<K,V>。这个RecordReader包装在Iterator[(K,V)]类内,复写Iterator的next()和hasNext方法,让compute返回的InterruptibleIterator[(K,V)]能够被迭代获得RecordReader取到的数据。

getPreferredLocations(split: Partition)操作:
  1. theSplit.serializableHadoopSplit.value.getLocations.filter(_ != "localhost")
复制代码
在NewHadoopPartition里SerializableWritable将split序列化,然后调用InputSplit本身的getLocations接口,得到有数据分布节点的nodes name列表。
WholeTextFileRDD
NewHadoopRDD的子类


  1. private[spark] class WholeTextFileRDD(
  2.     sc : SparkContext,
  3.     inputFormatClass: Class[_ <: WholeTextFileInputFormat],
  4.     keyClass: Class[String],
  5.     valueClass: Class[String],
  6.     @transient conf: Configuration,
  7.     minSplits: Int)
  8.   extends NewHadoopRDD[String, String](sc, inputFormatClass, keyClass, valueClass, conf) {
复制代码

复写了getPartitions方法:

NewHadoopRDD有自己的inputFormat实现类和recordReader实现类。在spark/input package下专门写了这两个类的实现。感觉是种参考。

InputFormat
WholeTextFileRDD在spark里实现了自己的inputFormat。读取的File以K,V的结构获取,K为path,V为整个file的content。


复写createRecordReader以使用WholeTextFileRecordReader
复写setMaxSplitSize方法,由于用户可以传入minSplits数目,计算平均大小(splits files总大小除以split数目)的时候就变了。
RecordReader
复写nextKeyValue方法,会读出指定path下的file的内容,生成new Text()给value,结果是String。如果文件正在被别的进行打开着,会返回false。否则把file内容读进value里。


使用场景
在SparkContext下提供wholeTextFile方法


  1. def wholeTextFiles(path: String, minSplits: Int = defaultMinSplits):
  2.   RDD[(String, String)]
复制代码
用于读取一个路径下的所有text文件,以K,V的形式返回,K为一个文件的path,V为文件内容。比较适合小文件。






已有(1)人评论

跳转到指定楼层
韩克拉玛寒 发表于 2014-10-18 09:03:38
很好,谢谢楼主辛苦了
回复

使用道具 举报

您需要登录后才可以回帖 登录 | 立即注册

本版积分规则

关闭

推荐上一条 /2 下一条