只讲一件事:你敲下一条 SQL,集群里到底发生了什么。
一、先记住一张“逻辑地图”
不管你用 YARN、K8s 还是 Standalone,本质只有 4 层:
Client 提交 ↓ Cluster Manager(资源调度) ↓ Spark Application ├─ Driver(控制中枢) └─ Executor(干活进程) ↓ Storage / Shuffle后面所有内容,都是在这张图里按时间线走。
假设你执行的是一条非常普通的 SQL:
SELECTdept_id,avg(salary)FROMempWHEREhire_date>='2020-01-01'GROUPBYdept_id;emp是 HDFS / S3 上的 Parquet 表- 有 200 个文件
一、提交过程
第 1 步:客户端只做“挂号”
你运行:
spark-submit\--masteryarn\--num-executors10\--executor-memory 4G\sql_job.pyspark-submit本身不执行任何计算,它只干三件事:
- 把你的代码、依赖、配置打包
- 向 Cluster Manager(这里是 YARN)申请资源
- 说一句话:
“帮我启动一个 Driver”
📌 这一步结束,任务还没开始跑,连 SQL 都没解析。
第 2 步:Driver 启动,大脑上线
YARN 分配一个 Container,启动Driver JVM。
Driver 里几个关键角色:
| 模块 | 干啥用 |
|---|---|
| SparkContext | 整个应用的入口 |
| DAGScheduler | 把 SQL / 代码变成 DAG |
| TaskScheduler | 把 Task 发给 Executor |
| SchedulerBackend | 和 YARN 沟通资源 |
⚠️ 重要认知:
- Driver不存业务数据
- Driver不算 salary、不算 avg
- 它只负责:解析、规划、调度、收状态
第 3 步:Executor 是“工人”
Driver 向 YARN 说:
“我要 10 个 Executor,每个 4G 内存”
为什么要 10 个?
这是你指定的资源配额,不是 Spark 算出来的。
- Executor = 进程
- 每个 Executor 里有很多 Task 线程
- 真正决定“同时能跑多少活”的是:
并行度 = Executor 数 × executor-cores例如:
10 Executor、每个 4 core → 最多 40 个 Task 同时跑
📌 Executor 数量 ≠ Task 数量,后面会看到 Task 远多于 10 个。
Executor 启动后,会向 Driver 注册:
“我上线了,可以接活。”
第 4 步:SQL 在 Driver 里被“拆”
1️⃣ SQL → 逻辑计划
Catalyst 把 SQL 解析成一棵树:
Aggregate [dept_id] Project [dept_id, salary] Filter (hire_date >= '2020-01-01') Scan Parquet2️⃣ 优化(只在 Driver 里改“计划”)
典型优化:
- 谓词下推:
WHERE hire_date >= 2020推到 Scan - 列裁剪:只读
dept_id, salary, hire_date - Parquet 列存裁剪
👉 这一步完全不碰数据,只是把“怎么读”定好。
3️⃣ 物理计划
变成 Spark 算子:
HashAggregate └─ HashAggregate └─ Scan Parquet并决定:读多少 Partition、 每个 Task 读哪一块文件
关键认知:Stage 为什么被切开?
DAGScheduler 一看计划:
GROUP BY dept_id→ 同一个 dept_id 必须凑到一起
但数据是分散的,怎么办?
👉必须 Shuffle
规则:
有 Shuffle,就切 Stage
于是 DAG 被切成两段:
Stage 0:Filter + Partial Aggregate(Map) | Shuffle | Stage 1:Final Aggregate(Reduce)第五步、Stage 0 在 Executor 里到底干了啥?
Map Task 从哪来?
表有 200 个 Parquet 文件 → Stage 0 有200 个 Map Task
Driver 把这 200 个 Task 分批发给 Executor。假设 Executor 1 拿到 Task 1、Task 2。
Task 内部执行流程:
读数据
- BlockManager 从 HDFS 读一个文件块
- 优先读本地节点(数据本地性)
Filter
- 过滤掉
hire_date < 2020的员工
- 过滤掉
Partial Aggregate
- 不急着算 avg,而是先算:
(dept_id, sum(salary), count) - 这是“局部汇总”
- 不急着算 avg,而是先算:
第六步、Shuffle:Spark 最“脏”的地方
Map 端写 Shuffle
每个 Map Task 不会只写一个文件,而是:
- 对
dept_id做 hash - 按 Reduce 分区写
比如默认:
spark.sql.shuffle.partitions = 200那么:
- 有 200 个 Reduce Task
- 每个 Map Task 写 200 个小数据段
Map Task 1: → Reduce 0: (dept=10, sum=18000, cnt=2) → Reduce 1: (dept=20, sum=9000, cnt=1) ...写的是Executor 本地磁盘。
Reduce 端怎么读?
Reduce Task 3:
去所有 Map Task那里,读“属于分区 3”的那一份
也就是:
Map Task 1 → 读它的 partition 3 Map Task 2 → 读它的 partition 3 ... Map Task 200 → 读它的 partition 3✅ 所以:
每个 Reduce Task 会拉取所有 Map Task 的一部分数据
- 通过网络(Netty)
- 拉到内存 → 溢写磁盘 → 排序 → 聚合
📌 Shuffle 数据:
- 不在 Driver
- 不在 HDFS
- 就在 Executor 的磁盘 + 网络里
第七步、Stage 1:Final Aggregate
Stage 1 是200 个 Reduce Task。
以某个 Reduce Task 为例:它拉到的是同一个dept_id的所有局部 sum / count:
sum = 18000 + 6000 + ... count = 2 + 1 + ... avg = sum / count算完后:
- 如果是
SELECT→ 结果被 Driver 收集,返回客户端 - 如果是
INSERT→ Executor 直接写 HDFS / 表
三、把“资源”和“计算”彻底分清
很多人混淆这两件事,一定要拆开:
| 概念 | 决定因素 |
|---|---|
| Executor 数 | num-executors(你配的) |
| 每个 Executor 能力 | executor-cores |
| Map Task 数 | 输入文件数 / Partition 数 |
| Reduce Task 数 | spark.sql.shuffle.partitions |
所以你看到的现象是:
- 10 个 Executor
- 但 Stage 0 有 200 个 Task
- Executor 轮流接 Task,跑完一个接下一个
四、用一句话串完整流程
你提交 SQL → YARN 启动 Driver → Driver 解析 SQL 成 DAG → 切出 Stage → 申请 Executor → Map Task 读文件、过滤、局部聚合 → 按 key 写 Shuffle → Reduce Task 跨节点拉数据 → 全局聚合 → 结果返回
三个最容易误解的点(记住就能秒杀面试)
Driver 不计算:它只调度、记状态、收心跳
Reduce Task 不是只拉一个 Map:它拉“所有 Map 里属于自己的那一块”
Executor ≠ Task
- Executor 是工人
- Task 是活
- 工人少,活可以很多,只是排队干