Spark 核心概念笔记
前面整理完 Hadoop 之后,自然要问一个问题:既然有 MapReduce 了,为什么还需要 Spark?
这篇笔记试图回答的就是这个。理清之后我发现,Spark 相比 MapReduce 的改进,核心其实只有一件事——把中间结果放在内存里。但由这一个改动带来的连锁反应,才是它真正的价值所在。
MapReduce 慢在哪
要先理解 Spark,得先知道 MapReduce 的问题。
MapReduce 的执行模型很死板:一个作业就是 Map 阶段和 Reduce 阶段,中间隔着一次 Shuffle。每个阶段的输出都要落盘——map 的输出写到本地磁盘,reduce 的结果写到 HDFS。
问题在于:
一、落盘是刚需但不是必然。写磁盘、读磁盘、序列化、反序列化,这些都是纯粹的开销。很多中间结果其实还可以继续用,但 MapReduce 用完就扔。
二、复杂任务要拆成多个作业串联。比如"先过滤、再分组统计、最后排序"这个需求,MapReduce 里得写成多个 Job 首尾相接,每个 Job 都要完整走一遍落盘流程。
结果就是:MapReduce 的瓶颈往往不在计算,而在 I/O。CPU 大部分时间在等磁盘。
Spark 的核心改进:RDD
Spark 的答案是把中间结果留在内存里。但它不是简单地"用内存代替磁盘",而是围绕一个新的抽象来设计的——RDD。
RDD 是什么
RDD 全称 Resilient Distributed Dataset(弹性分布式数据集)。拆开看每个词都有含义:
| 词 | 含义 |
|---|---|
| Resilient 弹性 | 某个分区的数据丢了,能根据血缘关系重新算出来,不需要存副本 |
| Distributed 分布式 | 数据分片存储在多台机器上,并行处理 |
| Dataset 数据集 | 是一个只读的元素集合 |
其中"弹性"这一条设计得很巧。HDFS 靠多存几份副本来容错,代价是存储开销翻几倍;RDD 靠记住它是怎么算出来的来容错——这个"怎么算出来的"叫血缘(lineage)。
举个例子:如果一个 RDD 是从某个 HDFS 文件 map 出来的,那么某个分区的数据丢了,Spark 只要对对应分区的原始数据重新执行一次 map,就能把它恢复出来。不需要存副本,也不需要人工介入。
这个设计的前提是:计算过程是确定性的、可重放的。
惰性求值:算不算,先记账
RDD 的操作分两类,这个概念是理解 Spark 执行流程的关键:
转换(Transformation)——map、filter、flatMap 这些。它们的特点是:调用时并不真的执行,只是记下"要对数据做什么"。
行动(Action)——count、collect、saveAsTextFile 这些。只有遇到行动操作,之前记下的所有转换才会真正执行。
这就是惰性求值(lazy evaluation)。
一开始我觉得这是多此一举——为什么不写一行执行一行?后来才明白它的价值:因为 Spark 能看到完整的计算链条,才有机会优化它。
举个例子,如果连续写了三个 filter,Spark 可以把它们合并成一个——反正数据只扫一遍就够了。如果边写边执行,就没有这个优化机会。
一个实用的推论:如果你写了一段 Spark 代码,跑起来什么都不发生,大概率是忘了加行动操作。转换只在记账,不记账到行动就不会触发。
依赖关系与 Stage 划分
行动操作触发执行时,Spark 会把整个计算链条画成一张 DAG(有向无环图),然后据此划分阶段。
划分的依据是依赖类型:
窄依赖 vs 宽依赖
| 类型 | 含义 | 典型操作 | 能否流水线 |
|---|---|---|---|
| 窄依赖 | 父 RDD 的每个分区,最多只被子 RDD 的一个分区使用 | map、filter、union | 能 |
| 宽依赖 | 父 RDD 的一个分区,会被子 RDD 的多个分区使用 | groupByKey、reduceByKey | 不能 |
为什么这个区分重要?因为在窄依赖下,每个分区的计算是独立的,可以在一个阶段内流水线跑完,不用等别的分区。
而宽依赖意味着数据要重新洗牌(shuffle)——父 RDD 的每个分区都要按规则拆开,分发给子 RDD 的各个分区。这必须等所有父分区都算完才能开始,所以它就成了阶段的分界线。
于是 Spark 的 Stage 划分规则可以一句话概括:从前往后走,遇到宽依赖就切一刀。
Stage 1(全是窄依赖,流水线执行)
map → filter → flatMap
↓ 遇到宽依赖(如 groupByKey),切分
Stage 2
reduceByKey → ...
↓
触发行动操作有个细节值得注意:每个 Stage 内部是一条流水线。数据在内存里依次流过 map、filter,不需要每个算子都生成一份完整的中间数据。这就是它比 MapReduce 快的重要原因之一——MapReduce 每个阶段都要落盘,阶段之间无法流水线。
与 MapReduce 的对比
把上面的内容整理成一张对照表:
| 维度 | MapReduce | Spark |
|---|---|---|
| 中间结果 | 落磁盘 | 尽量留内存 |
| 执行模型 | 固定的 Map → Reduce 两阶段 | DAG,任意算子组合 |
| 任务粒度 | 粗(一个 Job 一个阶段) | 细(Stage 内流水线) |
| 容错方式 | 中间结果落盘,可重跑 | 血缘关系,按需重算 |
| 迭代计算 | 每次都重新读磁盘 | 内存中反复迭代,快得多 |
| 适合场景 | 一次性的批处理 | 迭代计算、交互式查询、流处理 |
最后两行是重点。MapReduce 在"迭代计算"上的劣势尤其明显——机器学习训练、图计算这类任务要反复扫同一份数据几十上百轮,每一轮都从磁盘读一遍,速度差距能到十倍以上。这也是 Spark 最早吸引人的地方。
Spark 家族的组件
概念笔记最后把组件版图理一下:
| 组件 | 作用 |
|---|---|
| Spark Core | 最底层,提供 RDD 和执行引擎,其他组件都建立在它之上 |
| Spark SQL | 用 SQL 操作结构化数据,底层还是翻译成 RDD 操作 |
| Spark Streaming | 流处理,把流数据切成小批次处理 |
| MLlib | 机器学习库 |
| GraphX | 图计算 |
值得注意的一点:Spark 自己不负责存储。它通常跑在 HDFS 上读数据,也可能读 HBase、S3 或本地文件。它也不一定跑在 YARN 上——还有 Standalone 和 Mesos 两种资源管理模式。
所以严格说,Spark 不是"替代 Hadoop",而是"替代 MapReduce"。它和 HDFS、YARN 是互补关系,共同组成大数据的技术栈。
回过头看,这一章里我理解最深的是三个点。
Spark 快的根本原因不是"用了内存",而是"能看见完整计算链"。惰性求值让优化成为可能,DAG 让流水线和阶段划分成为可能——内存只是让这些优化真正发挥效果的条件。如果只记住"Spark 是用内存的",那遇到它不快的情况就会想不通。
宽依赖是性能的分水岭。能用窄依赖表达的逻辑就不要用宽依赖,因为宽依赖必然触发 shuffle,而 shuffle 是分布式计算里最贵的操作。写 Spark 代码时,"这个算子会不会引起 shuffle"应该成为一个条件反射。
RDD 用血缘换掉了副本。这是一次很聪明的设计取舍——用"重算"的代价换掉了"多存几份"的存储开销,前提是计算必须可重放。这也顺带解释了为什么 RDD 是只读的:如果允许随便改,血缘就记不清了,重算也就无从谈起。
最后补一句提醒:Spark 比 MapReduce 快,但这个"快多少"高度依赖场景。数据规模、shuffle 的数据量、中间结果能不能装进内存,任何一项变了,倍数都可能差出数量级。如果数据集本身很小,或者中间结果远超内存容量导致频繁落盘,Spark 的优势会被大幅削弱——它的设计前提是"中间结果放得下",前提不成立时,收益也就有限。
脱离具体场景谈"快 N 倍"没有意义。理解它为什么快,比记住一个倍数更有用。
暂无评论