欢迎光临
我们一直在努力

【 Spark 架构】一次 SQL 从提交到跑完的全景拆解

只讲一件事:你敲下一条 SQL,集群里到底发生了什么。

一、先记住一张“逻辑地图”

不管你用 YARN、K8s 还是 Standalone,本质只有 4 层:

Client 提交

Cluster Manager(资源调度)

Spark Application
├─ Driver(控制中枢)
└─ Executor(干活进程)

Storage / Shuffle

后面所有内容,都是在这张图里按时间线走。


假设你执行的是一条非常普通的 SQL:

SELECT dept_id, avg(salary)
FROM emp
WHERE hire_date >= '2020-01-01'
GROUP BY dept_id;

  • emp 是 HDFS / S3 上的 Parquet 表
  • 有 200 个文件

一、提交过程

第 1 步:客户端只做“挂号”

你运行:

spark-submit \\
–master yarn \\
–num-executors 10 \\
–executor-memory 4G \\
sql_job.py

spark-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 Parquet

    2️⃣ 优化(只在 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)
    • 这是“局部汇总”

  • 第六步、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 是活
    • 工人少,活可以很多,只是排队干
  • 赞(0)
    未经允许不得转载:171主机测评 » 【 Spark 架构】一次 SQL 从提交到跑完的全景拆解
    分享到: 更多 (0)

    评论 抢沙发

    • 昵称 (必填)
    • 邮箱 (必填)
    • 网址