只讲一件事你敲下一条 SQL集群里到底发生了什么。一、先记住一张“逻辑地图”不管你用 YARN、K8s 还是 Standalone本质只有 4 层Client 提交 ↓ Cluster Manager资源调度 ↓ Spark Application ├─ Driver控制中枢 └─ Executor干活进程 ↓ Storage / Shuffle后面所有内容都是在这张图里按时间线走。假设你执行的是一条非常普通的 SQLSELECTdept_id,avg(salary)FROMempWHEREhire_date2020-01-01GROUPBYdept_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 / 代码变成 DAGTaskScheduler把 Task 发给 ExecutorSchedulerBackend和 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_dateParquet 列存裁剪 这一步完全不碰数据只是把“怎么读”定好。3️⃣ 物理计划变成 Spark 算子HashAggregate └─ HashAggregate └─ Scan Parquet并决定读多少 Partition、 每个 Task 读哪一块文件关键认知Stage 为什么被切开DAGScheduler 一看计划GROUP BY dept_id→ 同一个 dept_id 必须凑到一起但数据是分散的怎么办必须 Shuffle规则有 Shuffle就切 Stage于是 DAG 被切成两段Stage 0Filter Partial AggregateMap | Shuffle | Stage 1Final AggregateReduce第五步、Stage 0 在 Executor 里到底干了啥Map Task 从哪来表有 200 个 Parquet 文件 → Stage 0 有200 个 Map TaskDriver 把这 200 个 Task 分批发给 Executor。假设 Executor 1 拿到 Task 1、Task 2。Task 内部执行流程读数据BlockManager 从 HDFS 读一个文件块优先读本地节点数据本地性Filter过滤掉hire_date 2020的员工Partial Aggregate不急着算 avg而是先算(dept_id, sum(salary), count)这是“局部汇总”第六步、ShuffleSpark 最“脏”的地方Map 端写 Shuffle每个 Map Task 不会只写一个文件而是对dept_id做 hash按 Reduce 分区写比如默认spark.sql.shuffle.partitions 200那么有 200 个 Reduce Task每个 Map Task 写 200 个小数据段Map Task 1: → Reduce 0: (dept10, sum18000, cnt2) → Reduce 1: (dept20, sum9000, cnt1) ...写的是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 1Final AggregateStage 1 是200 个 Reduce Task。以某个 Reduce Task 为例它拉到的是同一个dept_id的所有局部 sum / countsum 18000 6000 ... count 2 1 ... avg sum / count算完后如果是SELECT→ 结果被 Driver 收集返回客户端如果是INSERT→ Executor 直接写 HDFS / 表三、把“资源”和“计算”彻底分清很多人混淆这两件事一定要拆开概念决定因素Executor 数num-executors你配的每个 Executor 能力executor-coresMap Task 数输入文件数 / Partition 数Reduce Task 数spark.sql.shuffle.partitions所以你看到的现象是10 个 Executor但 Stage 0 有 200 个 TaskExecutor 轮流接 Task跑完一个接下一个四、用一句话串完整流程你提交 SQL → YARN 启动 Driver → Driver 解析 SQL 成 DAG → 切出 Stage → 申请 Executor → Map Task 读文件、过滤、局部聚合 → 按 key 写 Shuffle → Reduce Task 跨节点拉数据 → 全局聚合 → 结果返回三个最容易误解的点记住就能秒杀面试Driver 不计算它只调度、记状态、收心跳Reduce Task 不是只拉一个 Map它拉“所有 Map 里属于自己的那一块”Executor ≠ TaskExecutor 是工人Task 是活工人少活可以很多只是排队干