Spark 架构总览
概述
Spark 用内存计算把批处理性能提升了数倍到数十倍,其核心是DAG 执行引擎 + 内存计算。本文从整体架构入手:Driver/Executor 如何分工、Job/Stage/Task 三层模型、一个 Action 如何被拆成 DAG 执行,以及 SparkEnv 承载的运行时组件。
一、Spark 定位
| 对比 | MapReduce | Spark |
|---|---|---|
| 计算模型 | 每步落盘 | 内存 DAG |
| 中间结果 | 写 HDFS | 内存/磁盘缓存 |
| 延迟 | 分钟级 | 秒级(内存) |
| 编程模型 | Map/Reduce 二元 | RDD/DataFrame 丰富算子 |
| 迭代计算 | 极慢 | 快(数据常驻) |
Spark 的本质:把用户的算子序列转成有向无环图(DAG),尽量在内存中一次性算完,减少落盘。
二、运行时角色
2.1 角色分工
┌─────────────────────────────────────────────┐
│ Driver(提交节点) │
│ SparkContext │ DAGScheduler │ TaskScheduler │
│ 作业拆分、任务调度、结果收集 │
└──────────┬──────────────────────────────────┘
│ 提交 Task
┌───────┴────────┐
▼ ▼
┌────────┐ ┌────────┐ ┌────────┐
│Executor│ │Executor│ │Executor│
│(进程) │ │(进程) │ │(进程) │
│ Task 执行 │ │ Task 执行 │ │ Task 执行 │
│ 数据缓存 │ │ 数据缓存 │ │ 数据缓存 │
└────────┘ └────────┘ └────────┘
(由 ClusterManager 分配)| 角色 | 职责 |
|---|---|
| Driver | 运行 main 方法,持有 SparkContext,拆 DAG、调度 Task、收集结果 |
| Executor | 运行在 Worker 节点的进程,执行 Task、缓存 RDD 数据 |
| ClusterManager | 资源提供者:Standalone / YARN / K8s 中的 Master 或 RM |
| Worker/Node | 运行 Executor 的节点 |
2.2 各角色配置要点
| 项 | 参数示例 | 说明 |
|---|---|---|
| Driver 内存 | spark.driver.memory | 结果收集与调度开销 |
| Executor 内存 | spark.executor.memory | 执行与缓存 |
| Executor 核数 | spark.executor.cores | 并行 Task 数 |
| Executor 数量 | spark.executor.instances | 总 Executor 数 |
三、Job / Stage / Task 三层模型
3.1 三层概念
| 层级 | 触发 | 说明 |
|---|---|---|
| Job | 一个 Action(count/collect/save) | 一次完整计算 |
| Stage | DAG 按宽依赖切分 | 一组可并行的 Task |
| Task | Stage 内按分区拆 | 最小执行单元,一个分区一个 Task |
3.2 切分规则
一个 Job 的算子链(DAG)
窄依赖(map/filter):不打断,合并进同一 Stage
宽依赖(groupBy/join/reduceByKey):Shuffle 边界,切分新 Stagemap → filter → map (窄,同一 Stage)
│
▼
groupByKey (宽,Stage 边界)
│
▼
map → reduce(窄,新 Stage)Stage 数 = 宽依赖数 + 1。Stage 内任务完全并行,Stage 间顺序依赖。
四、DAG 执行流程
4.1 完整流程
Action(如 count)
│
1. SparkContext 提交 Job
│
2. DAGScheduler 构建 RDD 依赖 DAG
│
3. 按宽依赖切分 Stage,确定各 Stage 任务
│
4. TaskScheduler 把 Task 分发给 Executor
│
5. Executor 并行执行 Task,窄依赖管道式计算
│
6. 宽依赖触发 Shuffle,数据重分区后进入下一 Stage
│
7. 结果回传 Driver,Job 完成4.2 关键类
| 组件 | 职责 |
|---|---|
| DAGScheduler | 拆 Stage、建 DAG、处理 Stage 级失败重试 |
| TaskScheduler | 把 Task 调度到 Executor,管理任务队列与本地性 |
| SchedulerBackend | 与资源管理器交互,申请/释放 Executor |
4.3 失败重试
| 层级 | 重试机制 |
|---|---|
| Stage 失败 | 重新提交该 Stage 及其 Task |
| Task 失败 | 按 spark.task.maxFailures 重试 |
| Executor 丢失 | 重建 Executor,从血统重算丢失分区 |
五、SparkEnv
5.1 SparkEnv 是什么
每个 Driver/Executor 进程内的运行时环境容器,承载全部核心组件:
| 组件 | 职责 |
|---|---|
| MemoryManager | 统一内存管理(执行/存储) |
| BlockManager | 数据块缓存与传输 |
| ShuffleManager | Shuffle 读写管理 |
| RpcEnv | 进程间通信 |
| BroadcastManager | 广播变量 |
| SerializerManager | 序列化管理 |
5.2 BlockManager 细节
每个 Executor 的 BlockManager:
├── 内存块(MemoryStore)
├── 磁盘块(DiskStore)
└── 与其他 Executor 的 BlockManager 通信(获取远端块)数据本地性由此实现:Task 尽量调度到拥有所需数据块的 Executor。
六、执行模式与作业生命周期
6.1 三种运行模式
| 模式 | Driver 位置 | 适用 |
|---|---|---|
| Client | 提交节点 | 交互式、调试 |
| Cluster | 集群内 | 生产批任务 |
| Local | 本机 | 开发测试 |
6.2 生命周期
提交(spark-submit)→ 分配 Driver → 申请 Executor
→ 执行 Job → 结果回传 → 释放资源(SparkContext.stop)七、常见问题速查
| 问题 | 原因与处理 |
|---|---|
| Executor 数不足 | 资源未申请够,检查 yarn 队列与 spark.executor 配置 |
| Task 全部失败 | 数据或依赖问题,看 Driver 日志定位 Stage |
| 结果 collect 到 Driver 内存溢出 | 大数据量不要 collect,写文件/聚合后取 |
| 数据本地性差 | 块不在本地,等待或调整 spark.locality.wait |
| Executor 丢失频繁 | 内存/GC 问题,调大 Executor 内存或并行度 |