spark组件说明
Spark 核心组件与调度:提交一个任务,Driver、Executor、Stage、Task 都经历了什么
Spark 跑一个任务,涉及四个核心角色:
| 角色 | 比喻 | 做什么 |
|---|---|---|
| Driver | 项目经理 | 解析代码、拆任务、调度 |
| Executor | 一线工人 | 真正执行 Task |
| Cluster Manager | 人事部 | 分配资源(CPU、内存) |
| Task | 最小工作单元 | 一个分区的数据处理 |
一个 Spark 任务从提交到结束,就是这四个角色协同完成一件事:把用户代码拆成 Task,分配到 Executor 上并行执行
四个角色各管什么

Driver:任务的”大脑”
- 解析用户代码,生成 DAG(有向无环图)
- 把 DAG 切成 Stage,Stage 再切成 Task
- 向 Cluster Manager 申请 Executor
- 把 Task 分配到 Executor 上跑
- 监控任务状态,失败了重试
Executor:任务的”手脚”
- 跑 Task,执行具体的
map、filter、reduce操作 - 缓存 RDD 数据(内存/磁盘)
- 向 Driver 汇报执行进度
Cluster Manager:资源的”管家”
- 管理集群资源(CPU、内存)
- 启动和停止 Executor
- 支持 Standalone / YARN / K8s / Mesos
Task:最小的执行单元
- 一个 Task 处理一个 RDD 分区
- 所有 Task 加起来 = 一个 Stage 的全部计算
完整调度流程:提交一个 WordCount

第一步:提交应用,启动 Driver
用户执行 spark-submit,Cluster Manager 启动 Driver 进程,Driver 创建 SparkContext。
spark-submit --master yarn --class WordCount myapp.jar
第二步:Driver 申请 Executor
Driver 向 Cluster Manager 说:”我要 2 个 Executor,每个 2 核 4GB。”
Cluster Manager 在 Worker 节点上启动 Executor。Executor 启动后向 Driver 注册:”我准备好了,随时可以接活。”
第三步:Driver 生成执行计划
Driver 解析 WordCount 代码:
sc.textFile("hdfs://data.txt")
.flatMap(_.split(" "))
.map((_, 1))
.reduceByKey(_ + _)
.collect()
生成 DAG:
textFile → flatMap → map → reduceByKey → collect
按 Shuffle 边界切成两个 Stage:
Stage 1(无 Shuffle):textFile → flatMap → map → 输出 (word, 1)
Stage 2(有 Shuffle):reduceByKey(读 Stage 1 的数据)→ 聚合 → collect
第四步:Stage 拆成 Task
Stage 1 有 4 个分区 → 生成 4 个 Task
Stage 2 有 2 个分区(reduceByKey 的并行度)→ 生成 2 个 Task
第五步:Task 分配到 Executor
Driver 把 Task 发给 Executor。优先把 Task 分配到数据所在的节点(数据本地化)。
- 数据在 Executor 内存里 → PROCESS_LOCAL(最快)
- 数据在同一个节点的磁盘上 → NODE_LOCAL
- 数据在同一个机架 → RACK_LOCAL
- 数据在别的机架 → ANY(最慢)
第六步:Executor 执行 Task
Executor 用线程池并行跑 Task。跑完后把结果或中间数据存到 BlockManager(内存优先,不够写磁盘)。
第七步:结果收集,释放资源
collect() 把所有分区的结果拉到 Driver,合并后输出。
任务完成 → Driver 通知 Cluster Manager 关闭 Executor → 资源释放。
关键机制
1. Stage 划分规则:按 Shuffle 切
| 依赖类型 | 特点 | Stage 处理 |
|---|---|---|
| 窄依赖 | 父分区 → 子分区一对一 | 可以合并到同一个 Stage |
| 宽依赖 | 父分区 → 子分区多对多(需要 Shuffle) | 必须切 Stage |
map、flatMap、filter 是窄依赖 → 不用 ShufflereduceByKey、groupByKey、join 是宽依赖 → 触发 Shuffle,切 Stage
2. 数据本地化:Task 尽量靠近数据
优先级:PROCESS_LOCAL > NODE_LOCAL > RACK_LOCAL > ANY
如果数据在别的节点,Task 要跨网络拉数据,慢。Spark 会先等一会儿(spark.locality.wait),希望数据能本地化,等不到再降级。
3. 推测执行:慢 Task 开个”替补”
如果某个 Task 跑得特别慢(比其他 Task 慢 2 倍),Driver 会在另一个 Executor 上启动一个”替补”Task,谁先跑完用谁的结果。
配置:spark.speculation=true
4. 动态资源分配
YARN 模式下可以动态扩缩 Executor:
配置:spark.dynamicAllocation.enabled=true
任务忙的时候自动加 Executor,闲的时候自动释放。