第 300 题:实现一个简化版Spark的RDD,懒计算与血缘追踪。
题目
实现一个简化版Spark的RDD,懒计算与血缘追踪。
完整讲解
一、RDD 与 Spark 核心抽象
- RDD(Resilient Distributed Dataset):不可变、分区的分布式集合,支持 transformations(map、filter、groupBy 等,懒执行)与 actions(count、collect、save,触发计算)。懒计算即 transformation 只记录依赖与算子,不立即算,直到 action 触发;血缘(lineage) 即 RDD 的依赖图,用于故障时重算而无需持久化全部中间结果。
二、懒计算(Lazy Evaluation)
- 实现:每个 RDD 记录其 parent(s) 与 transformation 函数(如 map 的 f);调用 transformation 时只 new 一个 RDD 并保存 parent 与 f,不调用 f。Action 被调用时,从该 RDD 反向沿依赖递归到无依赖(或已有缓存),再正向按依赖顺序执行各 stage,得到结果。
- 好处:可做全局优化(如合并连续 map、推断分区)、按需计算、管道化与流水线。
三、血缘追踪(Lineage)
- 血缘:RDD 的依赖关系形成 DAG。窄依赖(父分区与子分区一对一或多对一)与宽依赖(shuffle,一对多)。记录方式:每个 RDD 存
dependencies: List[Dependency],如 OneToOneDependency、ShuffleDependency;每个 Dependency 指向父 RDD。
- 用途:(1)调度:按宽依赖划分 stage,stage 内 pipeline,stage 间 shuffle;(2)容错:某分区丢失时,从血缘追溯到源或缓存,只重算丢失分区及其依赖链;(3)持久化:用户可 persist 某 RDD,则重算时不必再溯到最前。
- 简化实现:RDD 基类含
dependencies()、compute(partition);子类如 MappedRDD 存 parent 与 f,compute 时对 parent 的该分区调 f。Action 如 count() 触发 sc.runJob(this, countPartition),调度器按血缘 DAG 提交 stage 并执行。
四、简化版要点
- 可不实现完整集群与 shuffle,只实现单机或伪分布;核心是「RDD = 依赖 + 算子」「transformation 懒、action 触发」「血缘 DAG 用于调度与重算」。可做 MapPartitionsRDD、FilterRDD、UnionRDD 等,以及 getDependencies、compute 的递归触发。
面试要点
- 能说清 RDD 的懒计算:transformation 只记依赖与函数,action 触发时才沿依赖计算。
- 能说明血缘的作用:记录依赖 DAG、用于调度(stage 划分)与容错(重算);能区分窄依赖与宽依赖。
记忆要点
- RDD:懒计算= transformation 不立刻算、action 触发;血缘=依赖 DAG,用于调度与故障重算。
- 实现:RDD 存 parent(s)+transformation;compute(partition) 递归父分区再算;stage 按宽依赖划分。