一、Spark 是什么
Apache Spark 是一个开源的统一分布式计算引擎,它以内存计算为核心,提供了 Java、Scala、Python、R 等高级 API,并提供以下功能:
- 批处理(Batch Processing)
- 交互式查询(Spark SQL)
- 流处理(Structured Streaming)
- 机器学习(MLlib)
- 图计算(GraphX)
1. Spark 生态全景图

-
Spark Core(核心基础):负责分布式任务调度、内存管理、容错、存储交互等基础功能。它引入了 RDD(弹性分布式数据集) 这一核心抽象,是所有其他组件的基石。
-
Spark SQL(结构化数据处理):用于处理结构化与半结构化数据,提供了DataFrame、Dataset,以及SQL查询引擎。底层也是基于RDD进行了封装优化。
-
Spark Streaming / Structured Streaming(流处理):基于 Spark SQL 的实时流数据处理引擎。
-
MLlib(机器学习库):分布式机器学习库,包含丰富的可扩展机器学习算法(分类、回归、聚类、推荐等),以及特征工程、模型评估等工作流工具,支持在分布式集群上训练大规模模型。
-
GraphX(图计算):用于图数据与图并行计算的 API 和算法库(如 PageRank、连通分量等),将图与集合操作统一,适合社交网络分析等场景。目前 Spark 生态中也推荐基于 DataFrame 的 GraphFrames 作为高级替代。
2. Spark 与 MapReduce 对比
| 计算模型 | 基于 DAG(有向无环图) 执行引擎,可将一个作业划分为多个阶段并行处理,并优化整体流程。 | 严格执行 Map → Shuffle → Reduce 的线性过程,一个复杂任务需要拆分为多个 MapReduce 作业串联,中间结果全程落盘。 |
| 速度 | 中间结果可缓存于内存,迭代计算(如机器学习)快 10–100 倍以上。 | 每次 Map/Reduce 结果都写入磁盘,导致大量 I/O 开销,速度较慢。 |
| 易用性 | 提供 DataFrame/Dataset 及丰富的算子(过滤、联接、聚合等),代码量远少于 MapReduce,并支持交互式查询。 | 编程相对底层,需要严格实现 Mapper 和 Reducer 类,逻辑复杂时代码冗长。 |
| 处理能力 | 原生支持流处理(Structured Streaming)、SQL、机器学习(MLlib)、图计算(GraphX),实现“一个堆栈解决多种需求”。 | 专为离线批处理设计,不支持实时处理,流处理需结合 Storm、Flink 等外部系统。 |
| 容错机制 | RDD 维护血统信息,分区丢失可基于依赖关系重算,同时配合检查点保证可靠性。 | 靠 Task 级别的重试机制,底层依赖于 HDFS 的副本实现容错。 |
| 资源管理 | Executor 进程常驻,可跨作业复用缓存数据,降低启动开销。 | 每个 Task 启动独立 JVM 进程,任务结束后即释放,无法共享内存数据。 |
二、核心数据抽象 RDD、DataFrame、Dataset
Spark 有三种数据表示方式:

1. RDD(Resilient Distributed Dataset 弹性分布式数据集)
Spark 最基础的分布式数据抽象,表示一个不可变、可分区的元素集合,支持并行操作。它的核心在于故障恢复:不通过数据复制,而是通过 血统(Lineage) 自动重算丢失的分区。弹性:可以基于内存存储也可以在磁盘中存储。
1.1 RDD 五大核心属性
每个 RDD 都具备以下 5 个内部属性:
-
分区列表(Partitions):是 Spark 对数据执行的逻辑分片,如果从 HDFS 读取,默认一个 HDFS Block(128MB)就是一个分区,可通过参数 spark.sql.files.maxPartitionBytes(默认128MB)控制单个分区最大读取的字节数。
-
表分区决定数据存在哪些目录和文件;HDFS Block 是这些文件的物理存储块;RDD 分区 是 Spark 在进行计算时的逻辑切分,读取阶段会尽量与 Block 对齐以获得数据本地性(一个 Block 通常完整存储在某个 DataNode 上,如果 Task 被调度到同一个节点,就可以直接读本地磁盘,避免网络传输),随后在 Shuffle 中重组。
- 并行度:每一个分区对应一个 Task(任务),所以有多少个分区,Spark 就可以同时起多少个 Task 并行计算。分区数通常是并行度的上限。分区划分Task。
-
表分区(目录) → 若干文件 → 每个文件按 HDFS Block 切分存储
↓
Spark 读取时:Block → RDD 分区(默认 1:1,受配置影响)
↓
再经过 Shuffle 等算子 → 新 RDD 分区(与 Block 无关)
-
计算函数(Compute):是对分区列表里每一个分区,具体执行数据生成或转换逻辑的代码块(函数),Spark 会把它分发到每个分区所在的 Executor 的 Task上执行。
-
你在代码里写的 rdd.map(…)、rdd.filter(…) 等算子,最终会层层组合成一个复合函数,由 compute 在分区上调用。
-
这个函数只在 行动算子 触发时才真正执行(惰性计算)。
-
-
依赖关系(Dependencies):记录父 RDD 与本 RDD 的关系(窄依赖/宽依赖),用于推断阶段划分和容错重算。
-
分区器(Partitioner,可选):专门用于键值对 RDD(PairRDD),它决定了数据在 Shuffle 阶段如何根据键分配到下游 RDD 的某个分区。只有需要根据 key 重分布数据(Shuffle)的操作才会有分区器。窄依赖没有分区器。常见的有:
-
HashPartitioner(默认):根据 Key 的 hashCode % 分区数 决定分区号。相同 key 必定同分区,但数据倾斜时可能个别分区过大。例如groupByKey、reduceByKey 等不要求排序的操作时。
-
RangePartitioner:对 key 进行采样预估数据分布情况,然后划分出大致相等的范围区间(数据分割边界),每个区间对应一个分区,最后将数据与分割边界进行比较放入对应的分区。要求 key 可排序,分区之间整体有序(分区内部不一定有序)。适用于需要全局排序或按范围划分的场景,例如 sortByKey。
-
-
首选位置(Preferred Locations,可选):返回每个分区数据所在的优选节点列表。Spark 调度器会尽可能把任务分配给这些节点。本地性优化:比如从 HDFS 读取时,首选位置就是该 Block 所在的数据节点,这样 Task 就可以直接读本地磁盘,省去网络传输。
1.2 计算函数的转换与行动算子 & 宽窄依赖关系
1.2.1 转换算子(Transformation):构建新的 RDD,形成依赖
转换算子从一个 RDD 生成另一个 RDD,它们只是定义了计算逻辑,不会立即执行。惰性求值(Lazy Evaluation):Transformation 只记录依赖关系,不触发计算;Action 才会提交 Job。
RDD 的血统(Lineage) 是指 RDD 之间通过转换操作形成的依赖关系链,记录了每一个 RDD 是如何从其父 RDD 一步步计算而来的“出身族谱”。
- 每个 RDD 对象内部都保存着对父 RDD 的引用以及转换函数。
- 这些引用串联起来,就形成了一张有向无环图(DAG)—— 即血统。
- 血统不仅记录了依赖关系,还记录了每个分区的数据来源(是来自父 RDD 的哪个分区,还是来自 Shuffle)。
Spark 的容错机制不靠数据复制,而靠血统重算:如果某个 RDD 的某个分区数据丢失(如 Executor 故障),Spark 不会全部重跑,而是根据血统向前回溯,找到最近的、可用或已持久化的父分区,重新执行 compute 函数链,只重算丢失的那几个分区。
Transformation 算子又可以分为窄依赖(无 Shuffle)和宽依赖(有 Shuffle)。

窄依赖:Narrow Dependency(pipeline 执行,无需跨分区拉取数据)
- 定义:父 RDD 的每个分区,最多被下游 RDD 的一个分区使用(不产生 Shuffle)。
- 管道执行:Task 直接在一个线程里顺序调用,无需网络传输。
- 分区列表不变:子 RDD 的分区列表与父 RDD 一一对应,Task 数量相同。
- 失败恢复高效:只需重算丢失的那个父分区即可。
| map | 对每个元素应用一个函数,返回新 RDD |
| flatMap | 类似 map,但一个输入可以映射为 0~多个输出(先映射再拍平) |
| filter | 保留满足条件的元素 |
| mapPartitions | 以分区为单位进行批量操作,常用于初始化连接 |
| sample | 按比例随机抽样 |
| union | 合并两个 RDD,不去重,分区简单相加 |
| 行过滤 | WHERE 条件 HAVING(仅过滤聚合结果) | 若 HAVING 过滤聚合后的数据,此时数据分区已固定,无需 Shuffle,属于窄依赖。 |
| 列投影与计算 | SELECT 字段(含 AS 别名) CAST(value AS type) CASE WHEN 分支判断 CONCAT / SUBSTRING 等内置函数(非聚合类) | 这些只对当前行做变换,不改变分区。 |
| 单表无排序的集合操作 | UNION ALL | 必须带 ALL。UNION(不带 ALL)底层会执行 DISTINCT,属于宽依赖。 |
| 采样 | TABLESAMPLE(x PERCENT) TABLESAMPLE(x ROWS) | Spark 会按物理块或种子随机采样,不改变分区依赖关系。 |
| 侧视图展开(Lateral View) | LATERAL VIEW EXPLODE(array_col) AS item | 一行展开为多行,但父子分区仍保持 1 对 1 映射(每个父分区数据只流入一个子分区),属于窄依赖。 |
| 局部排序(不全局排序) | SORT BY | SORT BY 只保证每个分区内部有序,不改变分区数量,是窄依赖。 DISTRIBUTE BY 是宽依赖(因为它会按 Key 重分布数据),千万不要记混! |
| 限制行数(取前N条) | LIMIT N | 物理上 Spark 会并行取每个分区的 Top N,再汇总,在第一个 Stage 内无需 Shuffle(后续如果有全局 ORDER BY 则另算)。 |
宽依赖:Wide Dependency / Shuffle Dependency(会触发 Shuffle,需要跨节点重新分区)
- 定义:父 RDD 的一个分区,会被多个子 RDD 分区使用(必须经过 Shuffle 重分区)。
- 需要 Shuffle 读写:父 Task 先按照分区器(如 HashPartitioner)将数据写入中间文件;子 Task 通过网络或本地读取属于自己分区的数据。
- 产生新的 Stage(阶段):宽依赖是 Stage 的边界,必须等上游 Stage 全部完成后,下游 Stage 才能启动。宽依赖划分Stage。
- 分区列表可能变化:子 RDD 的分区数由分区器(或 spark.sql.shuffle.partitions)决定,与父 RDD 无关。
| groupByKey | 按 Key 分组,返回 (K, Iterable<V>),不推荐,容易 OOM |
| reduceByKey | 按 Key 聚合,会在 Map 端做局部合并(Combiner),通常优于 groupByKey |
| sortByKey | 按 Key 对 RDD 进行全局排序 |
| join / leftOuterJoin / rightOuterJoin / fullOuterJoin | 连接两个 (K,V) 型的 RDD |
| cogroup | 对多个 RDD 按 Key 共同分组,join 的底层实现 |
| distinct | 去重(本质是先 map 成 (value, null),再 reduceByKey 最后 map 回来) |
| intersection / subtract | 交集 / 差集,会触发 Shuffle 去重 |
| repartition / coalesce | 调整分区数。repartition 一定会 Shuffle(可增可减);coalesce 默认窄依赖,可用于减少分区,但可设置 Shuffle |
| partitionBy | 只对 PairRDD 有效,按指定的分区器重新分区 |
-
聚合类:GROUP BY、DISTINCT、ROLLUP / CUBE / GROUPING SETS(这些需要按 Key 重新洗牌)
-
连接类:JOIN(包括 LEFT JOIN、INNER JOIN、RIGHT JOIN),除非被优化器转为 Broad




