spark组件说明

Spark 核心组件与调度:提交一个任务,Driver、Executor、Stage、Task 都经历了什么

Spark 跑一个任务,涉及四个核心角色:

角色 比喻 做什么
Driver 项目经理 解析代码、拆任务、调度
Executor 一线工人 真正执行 Task
Cluster Manager 人事部 分配资源(CPU、内存)
Task 最小工作单元 一个分区的数据处理

一个 Spark 任务从提交到结束,就是这四个角色协同完成一件事:把用户代码拆成 Task,分配到 Executor 上并行执行

四个角色各管什么

spark架构

Driver:任务的”大脑”

  • 解析用户代码,生成 DAG(有向无环图)
  • 把 DAG 切成 Stage,Stage 再切成 Task
  • 向 Cluster Manager 申请 Executor
  • 把 Task 分配到 Executor 上跑
  • 监控任务状态,失败了重试

Executor:任务的”手脚”

  • 跑 Task,执行具体的 mapfilterreduce 操作
  • 缓存 RDD 数据(内存/磁盘)
  • 向 Driver 汇报执行进度

Cluster Manager:资源的”管家”

  • 管理集群资源(CPU、内存)
  • 启动和停止 Executor
  • 支持 Standalone / YARN / K8s / Mesos

Task:最小的执行单元

  • 一个 Task 处理一个 RDD 分区
  • 所有 Task 加起来 = 一个 Stage 的全部计算

完整调度流程:提交一个 WordCount

spark调度分析

第一步:提交应用,启动 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

mapflatMapfilter 是窄依赖 → 不用 Shuffle
reduceByKeygroupByKeyjoin 是宽依赖 → 触发 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,闲的时候自动释放。