当前位置: 首页 > news >正文

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

只讲一件事:你敲下一条 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.py

spark-submit本身不执行任何计算,它只干三件事:

  1. 把你的代码、依赖、配置打包
  2. 向 Cluster Manager(这里是 YARN)申请资源
  3. 说一句话:

    “帮我启动一个 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 内部执行流程:

  1. 读数据

    • BlockManager 从 HDFS 读一个文件块
    • 优先读本地节点(数据本地性)
  2. Filter

    • 过滤掉hire_date < 2020的员工
  3. 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 跨节点拉数据 → 全局聚合 → 结果返回


三个最容易误解的点(记住就能秒杀面试)

  1. Driver 不计算:它只调度、记状态、收心跳

  2. Reduce Task 不是只拉一个 Map:它拉“所有 Map 里属于自己的那一块”

  3. Executor ≠ Task

    • Executor 是工人
    • Task 是活
    • 工人少,活可以很多,只是排队干
http://www.jsqmd.com/news/1402543/

相关文章:

  • HTTP到HTTPS强制跳转:301重定向与HSTS配置实战指南
  • OpenClaw v2026.3.11深度解析:AI智能体框架的安全、内核与跨平台进化
  • PyTorch+DeepSpeed大模型分布式训练实战指南
  • Batch Normalization 和 Layer Normalization 有什么用?
  • 音效素材资源网站推荐:2026 国内外正版与免费平台分类盘点
  • 2026年8月市面上可靠的智能制造能力成熟度评估公司推荐,CMMM,智能制造能力成熟度评估机构选哪家 - 企业权威推荐大使
  • WRC 2026前瞻:机器人产需共融下的核心能力与开发实践
  • Win11下VSCode配置Python虚拟环境:从venv原理到高效开发实战
  • 2026年8月山东铝合金铸造用金属硅/金属硅厂家推荐评估_山东鹏程光伏材料有限公司 - 行业平台推荐
  • Java学习笔记:Java流程控制
  • VICBench:多语言代码漏洞检测基准测试实战指南
  • 腾讯云轻量服务器安全加固指南:从SSH密钥到防火墙配置
  • 构建动态技能地图:从元技能到专业能力的系统化成长指南
  • Blender免费材质模型资源库全攻略:从PBR原理到高效搜索管理
  • HTTP状态码到底是个啥?一文看懂200、301、302、404、500
  • 构建智能工作流:开源流程引擎、专用小模型与智能体路由的集成实战
  • Commvault实战:Oracle数据库备份恢复全流程解析与避坑指南
  • MapInfo在线地图插件运行错误:Win10/Win11系统兼容性诊断与修复指南
  • 我的一点想法
  • Nacos 2.0 连接 127.0.0.1:9848 被拒绝?一文彻底解决 gRPC 端口通信问题
  • 2026年8月高强度微硅粉/微硅粉优质公司推荐_山东鹏程光伏材料有限公司 - 品牌宣传支持者
  • AI编程时代开发者核心竞争力:从编码到架构的升维竞争
  • 从词向量到语义搜索:Embedding原理与工程实践全解析
  • Ubuntu 22.04 Intel平台编译ALAMODE:从环境配置到性能调优完整指南
  • VSCode集成Cppcheck:Windows下C/C++代码静态分析与质量提升实战
  • Grok Bot插件生态解析:150+插件如何让AI从聊天机器人进化为可编程智能体
  • 大学生考哪些证书有用?2026年高含金量证书考证指南与就业避坑全解析
  • OpenAI API 集成实战:从环境配置到生产级代码助手开发
  • 2026年8月毛豆机采摘服务/跨区域毛豆机采摘服务农户推荐合作社_余姚康绿蔬菜专业合作社 - 行业平台推荐
  • FPGA/ASIC设计中set_input_delay约束详解:从原理到实战避坑指南