Spark 完全入门指南(上篇):基础与核心原理
面向有 IT 经验的工程师,不写废话,直接上干货。
如果你写过几年代码,懂分布式、懂数据库、懂 JVM,这篇会让你快速建立对 Spark 的深度认知,而不是停留在"会调 API"的层面。
目录
- 一、从 IT 老兵的视角看 Spark 的定位
- 二、Spark 架构与运行机制:不止是 Driver+Executor
- 三、RDD:分布式计算的最小抽象
- 四、宽窄依赖与 Stage 划分:Spark 调度的灵魂
- 五、内存管理机制:Spark 是怎么用内存的
- 六、调度机制:从 Job 到 Task 的完整链路
- 七、部署模式:Standalone / YARN / Mesos / K8s 深度对比
一、从 IT 老兵的视角看 Spark 的定位
1.1 大数据计算的三代演进
做了多年 IT 的人,应该能感受到数据处理技术的代际更替:
| 代际 | 代表技术 | 核心思想 | 解决的问题 | 局限 |
|---|---|---|---|---|
| 第一代 | 传统数据库(Oracle/MySQL) | 单机集中式计算 | 结构化数据、事务 | 数据量撑不住,扩展性差 |
| 第二代 | Hadoop MapReduce | 分布式批处理,移动计算不移动数据 | 海量数据离线处理 | 慢(磁盘 IO)、API 简陋、只支持批 |
| 第三代 | Spark / Flink | 内存计算、DAG 调度、统一引擎 | 速度、易用性、多场景 | 内存贵、流处理仍有优化空间 |
💡 关键洞察: 每一代技术不是完全替代上一代,而是解决上一代的核心痛点。 传统数据库没死,MapReduce 也还在跑,但新项目的首选已经是 Spark/Flink 了。 技术选型要看场景,不是越新越好。
1.2 Spark 到底解决了什么问题?
从工程师的角度,Spark 解决了三个核心痛点:
痛点一:MapReduce 太慢
MapReduce 把中间结果写磁盘,一个复杂的 ETL 可能要串十几个 MR Job,每个之间都要落盘+读盘。 Spark 把中间结果放内存,用 DAG 把多个操作串成流水线,中间不落盘。 性能提升 10~100 倍,这是最直接的价值。
痛点二:大数据技术栈太碎
以前做一个大数据项目,可能要同时用:
- MapReduce 做批处理
- Hive 做 SQL 查询
- Storm 做实时计算
- Mahout 做机器学习
- Giraph 做图计算
每个框架都要学、要运维、要互相导数据,痛苦不堪。
Spark 说:我一个框架全搞定。 Spark Core + Spark SQL + Streaming + MLlib + GraphX,一套 API、一套集群、一种运维。 这就是"统一计算引擎"的价值。
痛点三:开发效率低
MapReduce 写个 WordCount 要几十行代码,写个 Join 要上百行。 Spark 用 Scala/Python,WordCount 一行搞定,复杂逻辑也能简洁表达。 开发效率提升不是一点半点。
1.3 Spark 不是银弹:它的局限在哪?
客观看待,Spark 也有不擅长的地方:
| 场景 | Spark 的表现 | 更适合的技术 |
|---|---|---|
| 毫秒级流处理 | 微批模型,最低 100ms 级 | Flink(真流,毫秒级) |
| 点查/OLTP | 不擅长,延迟高 | HBase、Redis |
| OLAP 即席查询 | 可以,但不是最快 | ClickHouse、Doris、Impala |
| 深度学习 | MLlib 只支持传统 ML | TensorFlow、PyTorch |
| 图计算 | GraphX 维护不活跃 | Neo4j、JanusGraph、DGL |
| 小数据量(<10GB) | 杀鸡用牛刀,启动开销大 | Pandas、DuckDB |
🎯 结论: Spark 是大数据领域的"全能选手",但不是"单项冠军"。 批处理、SQL、机器学习是它的强项;毫秒级流、OLAP、深度学习不是它的主场。 知道一个技术的边界,比知道它能做什么更重要。
二、Spark 架构与运行机制:不止是 Driver+Executor
2.1 经典架构回顾
大部分教程讲到这里就结束了:
┌──────────────────┐
│ Driver Program │
│ (SparkContext) │
└────────┬─────────┘
│ 申请资源/分配任务
┌──────────────┼──────────────┐
▼ ▼ ▼
┌──────────┐ ┌──────────┐ ┌──────────┐
│ Worker 1 │ │ Worker 2 │ │ Worker N │
│ │ │ │ │ │
│ Executor │ │ Executor │ │ Executor │
│ (JVM) │ │ (JVM) │ │ (JVM) │
└──────────┘ └──────────┘ └──────────┘但对有经验的工程师来说,这远远不够。我们要深入到每个角色的内部。
2.2 Driver 内部:不只是"发号施令"
Driver 是整个应用的大脑,它内部运行着这些关键组件:
| 组件 | 职责 | 类比 |
|---|---|---|
| SparkContext | 入口,连接集群 | 公司前台 |
| DAGScheduler | 把 Job 拆成 Stage,处理宽依赖 | 项目总监,拆大项目为阶段 |
| TaskScheduler | 把 Task 分配给 Executor | 组长,分配具体任务 |
| SchedulerBackend | 和集群管理器通信,申请资源 | HR,招工人 |
| BlockManager | 管理缓存数据、Shuffle 数据 | 仓库管理员 |
| LiveListenerBus | 事件总线,Spark UI 靠它 | 公司公告栏 |
Driver 的工作流程
用户代码触发 Action
│
▼
SparkContext.runJob()
│
▼
DAGScheduler 接收,划分 Stage
│
▼
生成 TaskSet(每个 Stage 一组 Task)
│
▼
TaskScheduler 拿到 TaskSet
│
▼
SchedulerBackend 找空闲 Executor
│
▼
Task 序列化后发到 Executor 执行
│
▼
Executor 执行完汇报结果
│
▼
Driver 汇总,返回给用户💡 重要理解: Driver 不是只在启动时工作,它在整个应用运行期间都在持续调度、监控、汇总。 所以 Driver 不能挂,挂了整个应用就挂了。 生产环境要考虑 Driver 的高可用(比如 YARN 的 cluster 模式,Driver 跑在集群里由 YARN 监控)。
2.3 Executor 内部:不止是"干活的"
Executor 是一个 JVM 进程,内部结构:
┌─────────────────────────────────────┐
│ Executor (JVM) │
├─────────────────────────────────────┤
│ 线程池(ThreadPool) │
│ ┌────┐ ┌────┐ ┌────┐ ┌────┐ │
│ │Task│ │Task│ │Task│ │Task│ ... │ ← 每个核一个 Task 线程
│ └────┘ └────┘ └────┘ └────┘ │
├─────────────────────────────────────┤
│ 内存区域 │
│ ┌──────────┐ ┌──────────┐ │
│ │ 执行内存 │ │ 存储内存 │ │
│ │(Execution)│ │(Storage) │ │
│ └──────────┘ └──────────┘ │
├─────────────────────────────────────┤
│ BlockManager(数据块管理) │
├─────────────────────────────────────┤
│ Shuffle 客户端(读/写 Shuffle 数据)│
└─────────────────────────────────────┘Executor 的关键特性
- 多线程执行:一个 Executor 可以同时跑多个 Task(由
spark.executor.cores决定),共享同一个 JVM 和内存 - 进程级隔离:不同应用的 Executor 是不同的 JVM 进程,互不干扰
- 数据本地性:Executor 会优先处理自己节点上的数据(HDFS 本地块、缓存数据)
- 心跳机制:Executor 定期向 Driver 汇报心跳,超时没心跳就认为挂了
⚠️ 一个常见误区: 很多人以为"一个 Executor 跑一个 Task",不对。 一个 Executor 可以同时跑多个 Task(数量等于配置的 cores 数),它们共享内存。 这也是为什么 Executor 内存不能太小——多个 Task 同时跑,内存要够分。
2.4 集群管理器:Spark 不自己管资源
Spark 自己不做资源管理,而是把这事儿外包了。支持四种集群管理器:
| 集群管理器 | 出品方 | 特点 | 适用场景 |
|---|---|---|---|
| Standalone | Spark 自带 | 简单、独立、功能完整 | 测试、纯 Spark 集群 |
| YARN | Hadoop | 生态好、和 HDFS 天然集成、资源调度成熟 | 生产环境最常用 |
| Mesos | Apache | 细粒度资源分配、支持多种框架 | 有 Mesos 集群的环境 |
| Kubernetes | CNCF | 云原生、容器化、弹性伸缩 | 云原生环境、K8s 集群 |
🎯 生产环境怎么选?
- 已有 Hadoop 集群 → YARN(90% 的传统大数据公司选这个)
- 云原生/K8s 体系 → Kubernetes(互联网新公司、云厂商)
- 只是学习/测试 → Standalone(最简单)
- Mesos 现在用得越来越少了,不用重点学
三、RDD:分布式计算的最小抽象
3.1 RDD 的本质:不是"分布式数组"那么简单
很多教程说 RDD 是"分布式的不可变数据集",这没错,但不够深入。
从源码角度看,RDD 本质上是一个计算的抽象,它包含:
abstract class RDD[T] {
// 1. 分区列表:数据被分成了多少份
def getPartitions: Array[Partition]
// 2. 每个分区的计算函数:给定一个分区,怎么算出数据
def compute(split: Partition, context: TaskContext): Iterator[T]
// 3. 依赖列表:我依赖哪些父 RDD
def getDependencies: Seq[Dependency[_]]
// 4.(可选)分区器:键值对 RDD 怎么分区
def partitioner: Option[Partitioner]
// 5.(可选)首选位置:这个分区在哪台机器上算最好
def getPreferredLocations(split: Partition): Seq[String]
}💡 关键洞察: RDD 不存数据!RDD 存的是"怎么算出数据"的描述。 真正的数据在 compute() 执行的时候才产生,要么从父 RDD 算出来,要么从外部数据源读出来。 这就是为什么 RDD 是"懒执行"的——它只是一个计算蓝图,不是数据本身。
3.2 RDD 的五大特性(源码级理解)
特性一:一系列分区(Partitions)
- 分区是 RDD 的基本组成单位,也是并行计算的基本单位
- 一个分区对应一个 Task
- 分区数 = Task 数 = 并行度
分区数怎么决定?
- 从 HDFS 读:默认等于文件的 Block 数(一个 Block 一个分区)
parallelize:默认等于集群总核数repartition(n)/coalesce(n):手动指定- Shuffle 后:默认
spark.sql.shuffle.partitions(200)
特性二:每个分区有计算函数
compute()方法定义了"给定一个分区,怎么产出数据"- 不同类型的 RDD 有不同的 compute 实现
- HadoopRDD:从 HDFS 读数据
- MapPartitionsRDD:对父 RDD 的每个分区执行 map 函数
- ShuffledRDD:拉取 Shuffle 数据
特性三:依赖关系(Dependencies)
- RDD 记录了自己依赖哪些父 RDD,以及依赖类型(宽/窄)
- 这就是 Lineage(血缘),容错的基础
特性四:分区器(Partitioner)
- 只有键值对 RDD(RDD[(K, V)])才有分区器
- 决定 key 被分到哪个分区
- 内置两种:HashPartitioner、RangePartitioner
- 可以自定义 Partitioner
特性五:首选位置(Preferred Locations)
- 告诉调度器:这个分区在哪台机器上算最好(数据本地性)
- HadoopRDD 会返回 HDFS Block 所在的节点
- 调度器会尽量把 Task 调度到数据所在的节点,减少网络传输
🎯 面试深度题:RDD 为什么是不可变的? 答:不可变带来几个好处:
- 容错简单:不用维护数据修改历史,根据血缘重算就行
- 线程安全:多线程并发访问不需要锁
- 易于缓存:不可变的数据可以放心缓存,不用担心被改
- 函数式编程:符合纯函数式的设计,转换操作返回新 RDD
代价是:每次转换都生成新 RDD 对象,有一定对象创建开销(但 Spark 有优化)。
3.3 转换算子 vs 行动算子:懒执行的设计哲学
为什么要懒执行?
从工程设计角度,懒执行(Lazy Evaluation)有三大好处:
好处一:可以做全局优化
如果每一步都立刻执行,Spark 就看不到完整的计算图,没法优化。 攒到一起执行,Spark 可以:
- 把多个 map 合并成一个(map 融合)
- 调整操作顺序,减少 Shuffle
- 谓词下推(先过滤再 join)
- 生成最优的 Stage 划分
这和数据库的查询优化器是一个思路。
好处二:避免不必要的计算
如果你只需要结果的前 10 条,Spark 可以只算够 10 条就停,不用全量计算。 如果立刻执行,就浪费了。
好处三:资源利用率高
攒到一起执行,可以更合理地分配资源,避免频繁申请释放。
常见算子分类
转换算子(Transformation)—— 懒执行:
| 类别 | 算子 | 宽/窄依赖 |
|---|---|---|
| 元素级 | map、flatMap、filter、mapPartitions | 窄 |
| 集合级 | distinct、union、intersection、subtract | distinct 是宽,union 是窄 |
| 键值对 | reduceByKey、groupByKey、aggregateByKey、combineByKey | 宽 |
| 排序 | sortBy、sortByKey | 宽 |
| 重分区 | repartition、coalesce(shuffle=true) | 宽 |
| 连接 | join、leftOuterJoin、cogroup | 宽 |
行动算子(Action)—— 触发执行:
| 类别 | 算子 | 说明 |
|---|---|---|
| 收集 | collect、take、first、takeSample | 把数据拉到 Driver |
| 统计 | count、countByKey、reduce、fold、aggregate | 聚合统计 |
| 遍历 | foreach、foreachPartition | 对每条数据执行操作 |
| 保存 | saveAsTextFile、saveAsObjectFile、saveAsSequenceFile | 写出到外部存储 |
⚠️ 重要提醒(生产环境踩坑无数):
collect()会把所有数据拉到 Driver 端,数据量大时 Driver 必然 OOM。 生产环境绝对不要对大数据集用 collect()。 要看数据用take(n),要保存用saveAsTextFile(),要遍历用foreachPartition()。
3.4 RDD 持久化:cache vs persist
为什么需要持久化?
默认情况下,RDD 是不存数据的,每次行动算子触发都会从头算一遍。 如果一个 RDD 被多次使用,每次都重算就很浪费。
持久化就是把 RDD 的计算结果存起来(内存/磁盘),下次用直接读,不用重算。
cache vs persist 的区别
// cache() 其实就是 persist(MEMORY_ONLY) 的简写
def cache(): this.type = persist(StorageLevel.MEMORY_ONLY)| 方法 | 存储级别 | 说明 |
|---|---|---|
cache() | MEMORY_ONLY | 只存内存,放不下的分区就不缓存,下次重算 |
persist(MEMORY_ONLY) | MEMORY_ONLY | 同上 |
persist(MEMORY_AND_DISK) | 内存+磁盘 | 内存放不下的溢写磁盘,推荐 |
persist(MEMORY_ONLY_SER) | 内存(序列化) | 序列化后更省内存,但要反序列化开销 |
persist(DISK_ONLY) | 只存磁盘 | 内存不够时用 |
persist(MEMORY_AND_DISK_2) | 内存+磁盘,存2份 | 容错性更好,但占空间 |
什么时候该缓存?
- ✅ 同一个 RDD 被多次使用(比如迭代计算、多次查询)
- ✅ 计算这个 RDD 很耗时(比如经过很多转换)
- ✅ 数据量适中,放得下内存
- ❌ 只用一次的 RDD(缓存了反而多一次写的开销)
- ❌ 数据量特别大,内存放不下(缓存了也会被淘汰)
缓存的释放
rdd.unpersist() // 手动释放缓存缓存不用了记得释放,不然占着内存。当然 Spark 也会用 LRU 自动淘汰,但手动释放更可控。
💡 和 checkpoint 的区别:
- cache/persist:存在 Executor 的内存/磁盘,应用结束就没了,血缘还在
- checkpoint:存在 HDFS 等可靠存储,应用结束还在,会截断血缘(容错更彻底)
- 迭代计算(比如机器学习)常用 checkpoint 防止血缘链太长
四、宽窄依赖与 Stage 划分:Spark 调度的灵魂
4.1 宽窄依赖的本质区别
这是 Spark 最核心的概念,没有之一。理解了这个,才算真懂 Spark。
窄依赖(Narrow Dependency)
定义:父 RDD 的每个分区,最多被子 RDD 的一个分区使用。
父RDD分区: [P1] [P2] [P3] [P4]
│ │ │ │
▼ ▼ ▼ ▼
子RDD分区: [C1] [C2] [C3] [C4]一对一,或者多对一(多个父分区对应一个子分区,但每个父分区只去一个地方)。
特点:
- 不需要 Shuffle,数据不用跨节点传输
- 可以流水线执行(一个分区的数据一口气算完所有窄依赖操作)
- 容错成本低(重算一个分区就行)
常见算子:map、filter、flatMap、union、mapPartitions、coalesce(不 shuffle)
宽依赖(Wide Dependency / Shuffle Dependency)
定义:父 RDD 的每个分区,可能被子 RDD 的多个分区使用。
父RDD分区: [P1] [P2] [P3]
│╲ │╱ │
│ ╲ │ ╱ │
▼ ▼▼ ▼ ▼
子RDD分区: [C1] [C2] [C3]一个父分区的数据要分发到多个子分区,必须 Shuffle。
特点:
- 必须 Shuffle,数据要跨节点传输,还要写磁盘
- 不能流水线,必须等所有父分区都算完才能开始子分区
- 容错成本高(一个子分区丢了,可能要重算多个父分区)
常见算子:reduceByKey、groupByKey、join、distinct、sortBy、repartition
🎯 快速判断技巧: 问自己:这个操作需不需要"把相同 key 的数据放到一起"?
- 需要 → 宽依赖(要 Shuffle)
- 不需要 → 窄依赖
4.2 Stage 划分算法
Spark 怎么把一个 DAG 拆成多个 Stage?
算法:从最后一个 RDD(触发 Action 的那个)开始,从后往前遍历依赖链:
- 遇到窄依赖 → 划到同一个 Stage
- 遇到宽依赖 → 在这里切一刀,开一个新 Stage
原始 DAG(从左到右是数据流向):
A --map--> B --filter--> C --reduceByKey--> D --map--> E --collect()
窄依赖 窄依赖 宽依赖 窄依赖 Action
从后往前划分 Stage:
Stage 1(最后): D --map--> E --collect()
↑ 这里遇到宽依赖,切一刀
Stage 0(前面): A --map--> B --filter--> C --reduceByKey为什么从后往前?
因为 Action 是计算的起点(触发点),从触发点往前追溯,才能知道哪些操作是"必须连续执行"的。
Stage 的类型
- ShuffleMapStage:产出 Shuffle 数据的 Stage(不是最后一个 Stage)
- ResultStage:最后一个 Stage,执行 Action 操作
4.3 Shuffle 过程深入
宽依赖触发 Shuffle,Shuffle 是 Spark 性能的第一杀手。深入理解 Shuffle 过程,是调优的基础。
Shuffle 分两个阶段
阶段一:Shuffle Write(Map 端)
- 每个 Map Task(ShuffleMapStage 里的 Task)把自己的输出按 key 分区
- 写到本地磁盘的 Shuffle 文件中
- 每个 Map Task 为每个 Reduce 分区写一个文件段(Segment)
Map Task 1 的输出:
┌─────────────────────────────────┐
│ 分区0 │ 分区1 │ 分区2 │ 分区3 │ ← 按 key hash 分到不同分区
└─────────────────────────────────┘
│ │ │ │
▼ ▼ ▼ ▼
写到本地磁盘的 Shuffle 文件阶段二:Shuffle Read(Reduce 端)
- 每个 Reduce Task 从所有 Map Task 那里拉取属于自己的那个分区的数据
- 拉过来后进行聚合(如果需要),然后继续后续计算
Reduce Task 0 需要拉取:
- Map Task 1 的分区0 数据
- Map Task 2 的分区0 数据
- Map Task 3 的分区0 数据
...
然后合并、聚合Shuffle 为什么慢?
- 磁盘 IO:Shuffle Write 要写磁盘,Shuffle Read 要读磁盘
- 网络 IO:数据要跨节点传输
- 序列化/反序列化:数据要序列化才能传
- 排序/聚合:Reduce 端可能要排序、聚合
- 等待:Reduce 必须等所有 Map 都完成才能开始
🎯 调优启示: 减少 Shuffle = 提升性能。能不 Shuffle 就不 Shuffle,能少 Shuffle 就少 Shuffle。 具体方法后面调优篇详细讲。
Shuffle 管理器的演进
| 版本 | Shuffle 管理器 | 特点 |
|---|---|---|
| Spark 0.x | HashShuffleManager | 每个 Map Task 为每个 Reduce 写一个文件,文件数爆炸 |
| Spark 1.1+ | HashShuffle(文件合并) | 引入文件合并,减少文件数 |
| Spark 1.2+ | SortShuffleManager | 按 key 排序后写一个文件,默认 |
| Spark 2.0+ | 只有 SortShuffleManager | HashShuffle 被移除 |
SortShuffleManager 的优点:
- 每个 Map Task 只写一个数据文件 + 一个索引文件
- 文件数少,对操作系统友好
- 支持排序(有些算子需要排序)
五、内存管理机制:Spark 是怎么用内存的
5.1 内存区域划分
Executor 的内存(spark.executor.memory)不是一锅粥,而是分成了几个区域:
┌──────────────────────────────────────────────┐
│ Executor 总内存 │
├──────────────────────┬───────────────────────┤
│ 预留内存(300MB) │ 可用内存 │
│ (Reserved Memory) │ (Usable Memory) │
│ 系统用,不可配置 │ = 总内存 - 300MB │
├──────────────────────┼───────────────────────┤
│ │ ┌─────────┬─────────┐│
│ │ │ 执行内存 │ 存储内存 ││
│ │ │Execution│ Storage ││
│ │ │ (shuffle │ (cache/ ││
│ │ │ /join/ │ persist││
│ │ │ sort) │ /shuffle││
│ │ │ │ 中间结果)││
│ │ └─────────┴─────────┘│
│ │ 由 spark.memory.fraction │
│ │ 和 spark.memory.storageFraction 控制 │
└──────────────────────┴───────────────────────┘各区域说明
| 区域 | 大小 | 用途 |
|---|---|---|
| 预留内存 | 固定 300MB | 系统内部使用,不可配置 |
| 用户内存(User Memory) | 可用内存 × (1 - fraction) | 用户代码中的数据结构、UDF 中的对象 |
| 执行内存(Execution Memory) | 可用内存 × fraction × (1 - storageFraction) | Shuffle、Join、Sort、聚合等计算的临时数据 |
| 存储内存(Storage Memory) | 可用内存 × fraction × storageFraction | cache/persist 的数据、广播变量 |
关键参数
spark.memory.fraction:执行+存储内存占可用内存的比例,默认 0.6spark.memory.storageFraction:存储内存占执行+存储总内存的比例,默认 0.5
💡 举个例子: Executor 内存 8GB(8192MB)
- 预留:300MB
- 可用:8192 - 300 = 7892MB
- 执行+存储:7892 × 0.6 = 4735MB
- 存储内存:4735 × 0.5 = 2367MB
- 执行内存:4735 × 0.5 = 2367MB
- 用户内存:7892 × 0.4 = 3157MB
5.2 统一内存管理:执行和存储可以互相借用
Spark 1.6 之后用的是统一内存管理(Unified Memory Management)。
核心规则:
- 执行内存可以借用存储内存:如果存储内存没用满,执行内存可以借用
- 存储内存可以借用执行内存:如果执行内存没用满,存储内存可以借用
- 执行内存优先级更高:如果执行内存需要用了,可以把存储内存缓存的数据淘汰掉(如果是 MEMORY_AND_DISK 就溢写磁盘)
- 存储内存不能抢占执行内存:存储借了执行的内存,执行要用的时候,存储必须还
执行内存和存储内存的动态关系:
初始分配:
┌──────────────┬──────────────┐
│ 执行 50% │ 存储 50% │
└──────────────┴──────────────┘
存储没用满,执行借用:
┌───────────────────┬───────────┐
│ 执行 70%(借了) │ 存储 30% │
└───────────────────┴───────────┘
执行需要更多,抢占存储:
┌────────────────────────┬──────┐
│ 执行 85%(抢占存储) │ 15% │ ← 存储的数据被淘汰/溢写
└────────────────────────┴──────┘🎯 设计意图: 动态调整,提高内存利用率。 有的作业计算多(执行内存需求大),有的作业缓存多(存储内存需求大), 固定比例会浪费,动态借用更灵活。
5.3 常见内存问题与排查
问题一:Executor OOM
可能原因:
- 数据倾斜,某个分区数据特别大
- 执行内存不够(Shuffle 数据量大)
- 用户内存不够(UDF 里创建了大对象)
- 缓存数据太多,占了太多内存
排查方法:
- 看 Spark UI → Executors 页,看哪个 Executor 挂了
- 看 Executor 日志,找 OOM 堆栈
- 看 Stage 详情,是不是数据倾斜(某个 Task 数据量特别大)
解决方法:
- 增加分区数,减小每个分区的数据量
- 处理数据倾斜
- 调大 Executor 内存
- 用 persist(MEMORY_AND_DISK),内存不够溢写磁盘
- 优化 UDF,不要创建大对象
问题二:GC 严重
现象:Spark UI 上 Task 的 GC Time 占比很高(超过 10%)
原因:创建了太多短生命周期对象,JVM 频繁 GC
解决方法:
- 用 Kryo 序列化(比 Java 序列化快、省内存)
- 数据序列化存储(MEMORY_ONLY_SER)
- 优化代码,减少对象创建
- 调大 Executor 内存(给 GC 更多喘息空间)
- 用 G1 GC(
spark.executor.extraJavaOptions=-XX:+UseG1GC)
六、调度机制:从 Job 到 Task 的完整链路
6.1 调度的层级
Spark 的调度分三层:
应用级别(Application)
│ 由集群管理器调度(YARN/K8s)
▼
Job 级别(Job)
│ 由 DAGScheduler 调度,拆成 Stage
▼
Stage 级别(Stage)
│ 由 TaskScheduler 调度,拆成 Task
▼
Task 级别(Task)
由 Executor 的线程池执行6.2 一个 Action 触发的完整调度流程
1. 用户代码调用 Action(如 collect())
│
▼
2. SparkContext.runJob() 提交 Job
│
▼
3. DAGScheduler 接收 Job
- 从最终 RDD 从后往前遍历依赖
- 遇到宽依赖就切分 Stage
- 生成 Stage 的依赖关系图(DAG of Stage)
- 没有父 Stage 的 Stage 先提交
│
▼
4. 每个 Stage 生成 TaskSet
- Task 数量 = RDD 分区数
- 每个 Task 处理一个分区
│
▼
5. TaskScheduler 接收 TaskSet
- 按调度策略(FIFO/FAIR)排队
- 为每个 Task 计算数据本地性
│
▼
6. SchedulerBackend 分配资源
- 找空闲的 Executor
- 考虑数据本地性(PROCESS_LOCAL > NODE_LOCAL > RACK_LOCAL > ANY)
│
▼
7. Task 序列化后发到 Executor
│
▼
8. Executor 反序列化 Task,在线程池中执行
│
▼
9. Task 执行完,结果回传给 Driver
│
▼
10. 所有 Task 完成 → Stage 完成 → 提交下一个 Stage
│
▼
11. 所有 Stage 完成 → Job 完成 → 返回结果给用户6.3 数据本地性(Data Locality)
Spark 会尽量把 Task 调度到数据所在的位置,减少网络传输。
本地性优先级从高到低:
| 级别 | 含义 | 说明 |
|---|---|---|
| PROCESS_LOCAL | 同一个进程 | 数据在当前 Executor 的内存里(缓存数据),最快 |
| NODE_LOCAL | 同一个节点 | 数据在当前节点的磁盘上(HDFS 本地块、Shuffle 数据) |
| RACK_LOCAL | 同一个机架 | 数据在同机架的其他节点上,要走交换机 |
| ANY | 任意位置 | 数据在任意地方,要跨机架传输,最慢 |
本地性等待机制
如果高本地性的节点暂时没有空闲资源,Spark 会等一会儿(spark.locality.wait,默认 3 秒),等不到就降级到下一级本地性。
💡 调优启示:
- 如果数据本地性很差(很多 Task 是 ANY 级别),说明数据分布和计算资源不匹配
- 可以考虑:让 Spark 节点和 HDFS 节点共部、增加副本数、调整等待时间
6.4 调度模式:FIFO vs FAIR
Spark 内部的 Task 调度有两种模式(spark.scheduler.mode):
| 模式 | 特点 | 适用场景 |
|---|---|---|
| FIFO(默认) | 先提交的 Job 优先执行,前面的占满资源后面的等 | 单用户、批处理 |
| FAIR | 公平调度,多个 Job 轮流使用资源,支持权重 | 多用户、多租户、交互式查询 |
FAIR 模式还支持配置调度池(Pool),给不同的池设置权重和最小资源,实现更精细的资源隔离。
🎯 生产环境建议:
- 纯批处理、单用户 → FIFO(简单、 overhead 小)
- 多用户共享集群、有交互式查询 → FAIR(避免大作业饿死小作业)
七、部署模式:Standalone / YARN / Mesos / K8s 深度对比
7.1 客户端模式 vs 集群模式
这是所有部署模式都有的一个概念:Driver 跑在哪?
| 模式 | Driver 位置 | 特点 | 适用场景 |
|---|---|---|---|
| client(客户端模式) | Driver 跑在提交任务的机器上 | 日志直接在控制台看,方便调试;提交机器不能关 | 开发调试、交互式 |
| cluster(集群模式) | Driver 跑在集群里(由集群管理器管理) | Driver 由集群管理,有高可用;日志要去集群看 | 生产环境 |
⚠️ 重要: client 模式下,如果你提交完任务就关了电脑/终端,Driver 就挂了,整个任务就失败了。 生产环境一定要用 cluster 模式。
7.2 四种集群管理器对比
Standalone(Spark 自带)
架构:Master + Worker,Spark 自己管理资源
优点:
- 简单,不依赖其他组件
- 安装部署容易
- 功能完整,支持所有 Spark 特性
缺点:
- 只能跑 Spark,不能跑其他框架
- 资源管理不如 YARN 成熟
- 没有多租户、队列等高级功能
适用:学习测试、纯 Spark 的小集群
YARN(Hadoop 生态)
架构:ResourceManager + NodeManager,Spark 作为 YARN 的应用运行
优点:
- 和 HDFS 天然集成(数据本地性好)
- 资源调度成熟,支持队列、多租户
- 生态完善,Hive、HBase、Flink 都能跑在 YARN 上
- 生产环境最成熟稳定
缺点:
- 依赖 Hadoop 生态,组件多
- 弹性伸缩能力弱
- 容器化支持不好
适用:传统大数据公司、已有 Hadoop 集群
Mesos
架构:Mesos Master + Agent,细粒度资源分配
优点:
- 细粒度资源分配(可以按 CPU 核分配,不像 YARN 按 Container)
- 支持多种框架(Spark、Marathon、Chronos...)
- 资源利用率高
缺点:
- 社区活跃度下降,用的人越来越少
- 国内资料少
- 学习成本高
适用:已有 Mesos 集群的环境(越来越少)
Kubernetes(云原生)
架构:Spark 作为 K8s 的 Pod 运行,Driver 和 Executor 都是 Pod
优点:
- 云原生,和 K8s 生态无缝集成
- 弹性伸缩好(可以根据负载动态增减 Executor)
- 容器化,环境一致性好
- 资源隔离更彻底
- 支持 GPU 调度
缺点:
- 相对较新,生产案例不如 YARN 多
- 需要 K8s 运维能力
- 和 HDFS 的数据本地性需要特殊配置
适用:云原生架构、互联网新公司、混合云
7.3 选型决策树
你们有 Hadoop 集群吗?
├── 有 → 用 YARN(最稳妥,90% 的传统公司选这个)
└── 没有
├── 有 K8s 集群吗?
│ ├── 有 → 用 Kubernetes(云原生方向)
│ └── 没有
│ ├── 只是学习/测试 → Standalone(最简单)
│ └── 要上生产 → 建议搭 YARN 或 K8s
└── 有 Mesos 集群吗?
└── 有 → 用 Mesos(但建议考虑迁移到 K8s)🎯 趋势判断: 短期(1~3年):YARN 仍然是生产环境的主流 长期(3~5年):Kubernetes 会越来越多,特别是新建的集群 Mesos 会逐渐退出历史舞台 Standalone 主要用于学习和测试
上篇小结
到这里,你应该对 Spark 的核心原理有了深度理解:
- 定位:Spark 是第三代大数据计算引擎,解决了 MapReduce 慢、技术栈碎、开发效率低的问题,但不是银弹
- 架构:Driver + Executor 的主从架构,Driver 内部有 DAGScheduler、TaskScheduler 等组件
- RDD:本质是计算的抽象(不存数据,存计算描述),五大特性是理解 Spark 的钥匙
- 宽窄依赖:Spark 调度的灵魂,宽依赖触发 Shuffle,是性能瓶颈
- 内存管理:统一内存管理,执行和存储动态借用
- 调度:Job → Stage → Task 三层调度,数据本地性是关键优化点
- 部署:YARN 是当前主流,K8s 是未来方向
中篇预告:Spark SQL 深入(Catalyst 优化器、Tungsten 执行引擎)、Structured Streaming、MLlib、GraphX,以及 Spark vs Flink vs MapReduce 的深度对比。
文档版本:V1.0(上篇) 最后更新:2026 年 8 月 适用 Spark 版本:3.5.x