Skip to content

Spark 完全入门指南(上篇):基础与核心原理 ​

面向有 IT 经验的工程师,不写废话,直接上干货。

如果你写过几年代码,懂分布式、懂数据库、懂 JVM,这篇会让你快速建立对 Spark 的深度认知,而不是停留在"会调 API"的层面。


目录 ​


一、从 IT 老兵的视角看 Spark 的定位 ​

1.1 大数据计算的三代演进 ​

做了多年 IT 的人,应该能感受到数据处理技术的代际更替:

代际代表技术核心思想解决的问题局限
第一代传统数据库(Oracle/MySQL)单机集中式计算结构化数据、事务数据量撑不住,扩展性差
第二代Hadoop MapReduce分布式批处理,移动计算不移动数据海量数据离线处理慢(磁盘 IO)、API 简陋、只支持批
第三代Spark / Flink内存计算、DAG 调度、统一引擎速度、易用性、多场景内存贵、流处理仍有优化空间

💡 关键洞察: 每一代技术不是完全替代上一代,而是解决上一代的核心痛点。 传统数据库没死,MapReduce 也还在跑,但新项目的首选已经是 Spark/Flink 了。 技术选型要看场景,不是越新越好。

1.2 Spark 到底解决了什么问题? ​

从工程师的角度,Spark 解决了三个核心痛点:

痛点一:MapReduce 太慢 ​

MapReduce 把中间结果写磁盘,一个复杂的 ETL 可能要串十几个 MR Job,每个之间都要落盘+读盘。 Spark 把中间结果放内存,用 DAG 把多个操作串成流水线,中间不落盘。 性能提升 10~100 倍,这是最直接的价值。

痛点二:大数据技术栈太碎 ​

以前做一个大数据项目,可能要同时用:

  • MapReduce 做批处理
  • Hive 做 SQL 查询
  • Storm 做实时计算
  • Mahout 做机器学习
  • Giraph 做图计算

每个框架都要学、要运维、要互相导数据,痛苦不堪。

Spark 说:我一个框架全搞定。 Spark Core + Spark SQL + Streaming + MLlib + GraphX,一套 API、一套集群、一种运维。 这就是"统一计算引擎"的价值。

痛点三:开发效率低 ​

MapReduce 写个 WordCount 要几十行代码,写个 Join 要上百行。 Spark 用 Scala/Python,WordCount 一行搞定,复杂逻辑也能简洁表达。 开发效率提升不是一点半点。

1.3 Spark 不是银弹:它的局限在哪? ​

客观看待,Spark 也有不擅长的地方:

场景Spark 的表现更适合的技术
毫秒级流处理微批模型,最低 100ms 级Flink(真流,毫秒级)
点查/OLTP不擅长,延迟高HBase、Redis
OLAP 即席查询可以,但不是最快ClickHouse、Doris、Impala
深度学习MLlib 只支持传统 MLTensorFlow、PyTorch
图计算GraphX 维护不活跃Neo4j、JanusGraph、DGL
小数据量(<10GB)杀鸡用牛刀,启动开销大Pandas、DuckDB

🎯 结论: Spark 是大数据领域的"全能选手",但不是"单项冠军"。 批处理、SQL、机器学习是它的强项;毫秒级流、OLAP、深度学习不是它的主场。 知道一个技术的边界,比知道它能做什么更重要。


二、Spark 架构与运行机制:不止是 Driver+Executor ​

2.1 经典架构回顾 ​

大部分教程讲到这里就结束了:

                    ┌──────────────────┐
                    │   Driver Program │
                    │  (SparkContext)  │
                    └────────┬─────────┘
                             │ 申请资源/分配任务
              ┌──────────────┼──────────────┐
              ▼              ▼              ▼
        ┌──────────┐   ┌──────────┐   ┌──────────┐
        │ Worker 1 │   │ Worker 2 │   │ Worker N │
        │          │   │          │   │          │
        │ Executor │   │ Executor │   │ Executor │
        │  (JVM)   │   │  (JVM)   │   │  (JVM)   │
        └──────────┘   └──────────┘   └──────────┘

但对有经验的工程师来说,这远远不够。我们要深入到每个角色的内部。

2.2 Driver 内部:不只是"发号施令" ​

Driver 是整个应用的大脑,它内部运行着这些关键组件:

组件职责类比
SparkContext入口,连接集群公司前台
DAGScheduler把 Job 拆成 Stage,处理宽依赖项目总监,拆大项目为阶段
TaskScheduler把 Task 分配给 Executor组长,分配具体任务
SchedulerBackend和集群管理器通信,申请资源HR,招工人
BlockManager管理缓存数据、Shuffle 数据仓库管理员
LiveListenerBus事件总线,Spark UI 靠它公司公告栏

Driver 的工作流程 ​

用户代码触发 Action
        │
        ▼
  SparkContext.runJob()
        │
        ▼
  DAGScheduler 接收,划分 Stage
        │
        ▼
  生成 TaskSet(每个 Stage 一组 Task)
        │
        ▼
  TaskScheduler 拿到 TaskSet
        │
        ▼
  SchedulerBackend 找空闲 Executor
        │
        ▼
  Task 序列化后发到 Executor 执行
        │
        ▼
  Executor 执行完汇报结果
        │
        ▼
  Driver 汇总,返回给用户

💡 重要理解: Driver 不是只在启动时工作,它在整个应用运行期间都在持续调度、监控、汇总。 所以 Driver 不能挂,挂了整个应用就挂了。 生产环境要考虑 Driver 的高可用(比如 YARN 的 cluster 模式,Driver 跑在集群里由 YARN 监控)。

2.3 Executor 内部:不止是"干活的" ​

Executor 是一个 JVM 进程,内部结构:

┌─────────────────────────────────────┐
│            Executor (JVM)           │
├─────────────────────────────────────┤
│  线程池(ThreadPool)               │
│  ┌────┐ ┌────┐ ┌────┐ ┌────┐       │
│  │Task│ │Task│ │Task│ │Task│  ...  │  ← 每个核一个 Task 线程
│  └────┘ └────┘ └────┘ └────┘       │
├─────────────────────────────────────┤
│  内存区域                           │
│  ┌──────────┐ ┌──────────┐         │
│  │ 执行内存  │ │ 存储内存  │         │
│  │(Execution)│ │(Storage) │         │
│  └──────────┘ └──────────┘         │
├─────────────────────────────────────┤
│  BlockManager(数据块管理)          │
├─────────────────────────────────────┤
│  Shuffle 客户端(读/写 Shuffle 数据)│
└─────────────────────────────────────┘

Executor 的关键特性 ​

  1. 多线程执行:一个 Executor 可以同时跑多个 Task(由 spark.executor.cores 决定),共享同一个 JVM 和内存
  2. 进程级隔离:不同应用的 Executor 是不同的 JVM 进程,互不干扰
  3. 数据本地性:Executor 会优先处理自己节点上的数据(HDFS 本地块、缓存数据)
  4. 心跳机制:Executor 定期向 Driver 汇报心跳,超时没心跳就认为挂了

⚠️ 一个常见误区: 很多人以为"一个 Executor 跑一个 Task",不对。 一个 Executor 可以同时跑多个 Task(数量等于配置的 cores 数),它们共享内存。 这也是为什么 Executor 内存不能太小——多个 Task 同时跑,内存要够分。

2.4 集群管理器:Spark 不自己管资源 ​

Spark 自己不做资源管理,而是把这事儿外包了。支持四种集群管理器:

集群管理器出品方特点适用场景
StandaloneSpark 自带简单、独立、功能完整测试、纯 Spark 集群
YARNHadoop生态好、和 HDFS 天然集成、资源调度成熟生产环境最常用
MesosApache细粒度资源分配、支持多种框架有 Mesos 集群的环境
KubernetesCNCF云原生、容器化、弹性伸缩云原生环境、K8s 集群

🎯 生产环境怎么选?

  • 已有 Hadoop 集群 → YARN(90% 的传统大数据公司选这个)
  • 云原生/K8s 体系 → Kubernetes(互联网新公司、云厂商)
  • 只是学习/测试 → Standalone(最简单)
  • Mesos 现在用得越来越少了,不用重点学

三、RDD:分布式计算的最小抽象 ​

3.1 RDD 的本质:不是"分布式数组"那么简单 ​

很多教程说 RDD 是"分布式的不可变数据集",这没错,但不够深入。

从源码角度看,RDD 本质上是一个计算的抽象,它包含:

scala
abstract class RDD[T] {
  // 1. 分区列表:数据被分成了多少份
  def getPartitions: Array[Partition]

  // 2. 每个分区的计算函数:给定一个分区,怎么算出数据
  def compute(split: Partition, context: TaskContext): Iterator[T]

  // 3. 依赖列表:我依赖哪些父 RDD
  def getDependencies: Seq[Dependency[_]]

  // 4.(可选)分区器:键值对 RDD 怎么分区
  def partitioner: Option[Partitioner]

  // 5.(可选)首选位置:这个分区在哪台机器上算最好
  def getPreferredLocations(split: Partition): Seq[String]
}

💡 关键洞察: RDD 不存数据!RDD 存的是"怎么算出数据"的描述。 真正的数据在 compute() 执行的时候才产生,要么从父 RDD 算出来,要么从外部数据源读出来。 这就是为什么 RDD 是"懒执行"的——它只是一个计算蓝图,不是数据本身。

3.2 RDD 的五大特性(源码级理解) ​

特性一:一系列分区(Partitions) ​

  • 分区是 RDD 的基本组成单位,也是并行计算的基本单位
  • 一个分区对应一个 Task
  • 分区数 = Task 数 = 并行度

分区数怎么决定?

  • 从 HDFS 读:默认等于文件的 Block 数(一个 Block 一个分区)
  • parallelize:默认等于集群总核数
  • repartition(n) / coalesce(n):手动指定
  • Shuffle 后:默认 spark.sql.shuffle.partitions(200)

特性二:每个分区有计算函数 ​

  • compute() 方法定义了"给定一个分区,怎么产出数据"
  • 不同类型的 RDD 有不同的 compute 实现
    • HadoopRDD:从 HDFS 读数据
    • MapPartitionsRDD:对父 RDD 的每个分区执行 map 函数
    • ShuffledRDD:拉取 Shuffle 数据

特性三:依赖关系(Dependencies) ​

  • RDD 记录了自己依赖哪些父 RDD,以及依赖类型(宽/窄)
  • 这就是 Lineage(血缘),容错的基础

特性四:分区器(Partitioner) ​

  • 只有键值对 RDD(RDD[(K, V)])才有分区器
  • 决定 key 被分到哪个分区
  • 内置两种:HashPartitioner、RangePartitioner
  • 可以自定义 Partitioner

特性五:首选位置(Preferred Locations) ​

  • 告诉调度器:这个分区在哪台机器上算最好(数据本地性)
  • HadoopRDD 会返回 HDFS Block 所在的节点
  • 调度器会尽量把 Task 调度到数据所在的节点,减少网络传输

🎯 面试深度题:RDD 为什么是不可变的? 答:不可变带来几个好处:

  1. 容错简单:不用维护数据修改历史,根据血缘重算就行
  2. 线程安全:多线程并发访问不需要锁
  3. 易于缓存:不可变的数据可以放心缓存,不用担心被改
  4. 函数式编程:符合纯函数式的设计,转换操作返回新 RDD

代价是:每次转换都生成新 RDD 对象,有一定对象创建开销(但 Spark 有优化)。

3.3 转换算子 vs 行动算子:懒执行的设计哲学 ​

为什么要懒执行? ​

从工程设计角度,懒执行(Lazy Evaluation)有三大好处:

好处一:可以做全局优化 ​

如果每一步都立刻执行,Spark 就看不到完整的计算图,没法优化。 攒到一起执行,Spark 可以:

  • 把多个 map 合并成一个(map 融合)
  • 调整操作顺序,减少 Shuffle
  • 谓词下推(先过滤再 join)
  • 生成最优的 Stage 划分

这和数据库的查询优化器是一个思路。

好处二:避免不必要的计算 ​

如果你只需要结果的前 10 条,Spark 可以只算够 10 条就停,不用全量计算。 如果立刻执行,就浪费了。

好处三:资源利用率高 ​

攒到一起执行,可以更合理地分配资源,避免频繁申请释放。

常见算子分类 ​

转换算子(Transformation)—— 懒执行:

类别算子宽/窄依赖
元素级map、flatMap、filter、mapPartitions窄
集合级distinct、union、intersection、subtractdistinct 是宽,union 是窄
键值对reduceByKey、groupByKey、aggregateByKey、combineByKey宽
排序sortBy、sortByKey宽
重分区repartition、coalesce(shuffle=true)宽
连接join、leftOuterJoin、cogroup宽

行动算子(Action)—— 触发执行:

类别算子说明
收集collect、take、first、takeSample把数据拉到 Driver
统计count、countByKey、reduce、fold、aggregate聚合统计
遍历foreach、foreachPartition对每条数据执行操作
保存saveAsTextFile、saveAsObjectFile、saveAsSequenceFile写出到外部存储

⚠️ 重要提醒(生产环境踩坑无数): collect() 会把所有数据拉到 Driver 端,数据量大时 Driver 必然 OOM。 生产环境绝对不要对大数据集用 collect()。 要看数据用 take(n),要保存用 saveAsTextFile(),要遍历用 foreachPartition()。

3.4 RDD 持久化:cache vs persist ​

为什么需要持久化? ​

默认情况下,RDD 是不存数据的,每次行动算子触发都会从头算一遍。 如果一个 RDD 被多次使用,每次都重算就很浪费。

持久化就是把 RDD 的计算结果存起来(内存/磁盘),下次用直接读,不用重算。

cache vs persist 的区别 ​

scala
// cache() 其实就是 persist(MEMORY_ONLY) 的简写
def cache(): this.type = persist(StorageLevel.MEMORY_ONLY)
方法存储级别说明
cache()MEMORY_ONLY只存内存,放不下的分区就不缓存,下次重算
persist(MEMORY_ONLY)MEMORY_ONLY同上
persist(MEMORY_AND_DISK)内存+磁盘内存放不下的溢写磁盘,推荐
persist(MEMORY_ONLY_SER)内存(序列化)序列化后更省内存,但要反序列化开销
persist(DISK_ONLY)只存磁盘内存不够时用
persist(MEMORY_AND_DISK_2)内存+磁盘,存2份容错性更好,但占空间

什么时候该缓存? ​

  • ✅ 同一个 RDD 被多次使用(比如迭代计算、多次查询)
  • ✅ 计算这个 RDD 很耗时(比如经过很多转换)
  • ✅ 数据量适中,放得下内存
  • ❌ 只用一次的 RDD(缓存了反而多一次写的开销)
  • ❌ 数据量特别大,内存放不下(缓存了也会被淘汰)

缓存的释放 ​

scala
rdd.unpersist()  // 手动释放缓存

缓存不用了记得释放,不然占着内存。当然 Spark 也会用 LRU 自动淘汰,但手动释放更可控。

💡 和 checkpoint 的区别:

  • cache/persist:存在 Executor 的内存/磁盘,应用结束就没了,血缘还在
  • checkpoint:存在 HDFS 等可靠存储,应用结束还在,会截断血缘(容错更彻底)
  • 迭代计算(比如机器学习)常用 checkpoint 防止血缘链太长

四、宽窄依赖与 Stage 划分:Spark 调度的灵魂 ​

4.1 宽窄依赖的本质区别 ​

这是 Spark 最核心的概念,没有之一。理解了这个,才算真懂 Spark。

窄依赖(Narrow Dependency) ​

定义:父 RDD 的每个分区,最多被子 RDD 的一个分区使用。

父RDD分区:  [P1] [P2] [P3] [P4]
             │    │    │    │
             ▼    ▼    ▼    ▼
子RDD分区:  [C1] [C2] [C3] [C4]

一对一,或者多对一(多个父分区对应一个子分区,但每个父分区只去一个地方)。

特点:

  • 不需要 Shuffle,数据不用跨节点传输
  • 可以流水线执行(一个分区的数据一口气算完所有窄依赖操作)
  • 容错成本低(重算一个分区就行)

常见算子:map、filter、flatMap、union、mapPartitions、coalesce(不 shuffle)

宽依赖(Wide Dependency / Shuffle Dependency) ​

定义:父 RDD 的每个分区,可能被子 RDD 的多个分区使用。

父RDD分区:  [P1] [P2] [P3]
             │╲  │╱  │
             │ ╲ │ ╱ │
             ▼  ▼▼  ▼ ▼
子RDD分区:  [C1] [C2] [C3]

一个父分区的数据要分发到多个子分区,必须 Shuffle。

特点:

  • 必须 Shuffle,数据要跨节点传输,还要写磁盘
  • 不能流水线,必须等所有父分区都算完才能开始子分区
  • 容错成本高(一个子分区丢了,可能要重算多个父分区)

常见算子:reduceByKey、groupByKey、join、distinct、sortBy、repartition

🎯 快速判断技巧: 问自己:这个操作需不需要"把相同 key 的数据放到一起"?

  • 需要 → 宽依赖(要 Shuffle)
  • 不需要 → 窄依赖

4.2 Stage 划分算法 ​

Spark 怎么把一个 DAG 拆成多个 Stage?

算法:从最后一个 RDD(触发 Action 的那个)开始,从后往前遍历依赖链:

  • 遇到窄依赖 → 划到同一个 Stage
  • 遇到宽依赖 → 在这里切一刀,开一个新 Stage
原始 DAG(从左到右是数据流向):

A --map--> B --filter--> C --reduceByKey--> D --map--> E --collect()
            窄依赖        窄依赖      宽依赖       窄依赖     Action

从后往前划分 Stage:

Stage 1(最后):  D --map--> E --collect()
                 ↑ 这里遇到宽依赖,切一刀
Stage 0(前面):  A --map--> B --filter--> C --reduceByKey

为什么从后往前? ​

因为 Action 是计算的起点(触发点),从触发点往前追溯,才能知道哪些操作是"必须连续执行"的。

Stage 的类型 ​

  • ShuffleMapStage:产出 Shuffle 数据的 Stage(不是最后一个 Stage)
  • ResultStage:最后一个 Stage,执行 Action 操作

4.3 Shuffle 过程深入 ​

宽依赖触发 Shuffle,Shuffle 是 Spark 性能的第一杀手。深入理解 Shuffle 过程,是调优的基础。

Shuffle 分两个阶段 ​

阶段一:Shuffle Write(Map 端) ​

  • 每个 Map Task(ShuffleMapStage 里的 Task)把自己的输出按 key 分区
  • 写到本地磁盘的 Shuffle 文件中
  • 每个 Map Task 为每个 Reduce 分区写一个文件段(Segment)
Map Task 1 的输出:
  ┌─────────────────────────────────┐
  │ 分区0 │ 分区1 │ 分区2 │ 分区3   │  ← 按 key hash 分到不同分区
  └─────────────────────────────────┘
       │       │       │       │
       ▼       ▼       ▼       ▼
  写到本地磁盘的 Shuffle 文件

阶段二:Shuffle Read(Reduce 端) ​

  • 每个 Reduce Task 从所有 Map Task 那里拉取属于自己的那个分区的数据
  • 拉过来后进行聚合(如果需要),然后继续后续计算
Reduce Task 0 需要拉取:
  - Map Task 1 的分区0 数据
  - Map Task 2 的分区0 数据
  - Map Task 3 的分区0 数据
  ...
  然后合并、聚合

Shuffle 为什么慢? ​

  1. 磁盘 IO:Shuffle Write 要写磁盘,Shuffle Read 要读磁盘
  2. 网络 IO:数据要跨节点传输
  3. 序列化/反序列化:数据要序列化才能传
  4. 排序/聚合:Reduce 端可能要排序、聚合
  5. 等待:Reduce 必须等所有 Map 都完成才能开始

🎯 调优启示: 减少 Shuffle = 提升性能。能不 Shuffle 就不 Shuffle,能少 Shuffle 就少 Shuffle。 具体方法后面调优篇详细讲。

Shuffle 管理器的演进 ​

版本Shuffle 管理器特点
Spark 0.xHashShuffleManager每个 Map Task 为每个 Reduce 写一个文件,文件数爆炸
Spark 1.1+HashShuffle(文件合并)引入文件合并,减少文件数
Spark 1.2+SortShuffleManager按 key 排序后写一个文件,默认
Spark 2.0+只有 SortShuffleManagerHashShuffle 被移除

SortShuffleManager 的优点:

  • 每个 Map Task 只写一个数据文件 + 一个索引文件
  • 文件数少,对操作系统友好
  • 支持排序(有些算子需要排序)

五、内存管理机制:Spark 是怎么用内存的 ​

5.1 内存区域划分 ​

Executor 的内存(spark.executor.memory)不是一锅粥,而是分成了几个区域:

┌──────────────────────────────────────────────┐
│           Executor 总内存                     │
├──────────────────────┬───────────────────────┤
│  预留内存(300MB)    │  可用内存              │
│  (Reserved Memory)   │  (Usable Memory)      │
│  系统用,不可配置     │  = 总内存 - 300MB      │
├──────────────────────┼───────────────────────┤
│                      │  ┌─────────┬─────────┐│
│                      │  │ 执行内存 │ 存储内存 ││
│                      │  │Execution│ Storage ││
│                      │  │ (shuffle │ (cache/ ││
│                      │  │  /join/  │  persist││
│                      │  │  sort)   │  /shuffle││
│                      │  │         │  中间结果)││
│                      │  └─────────┴─────────┘│
│                      │  由 spark.memory.fraction │
│                      │  和 spark.memory.storageFraction 控制 │
└──────────────────────┴───────────────────────┘

各区域说明 ​

区域大小用途
预留内存固定 300MB系统内部使用,不可配置
用户内存(User Memory)可用内存 × (1 - fraction)用户代码中的数据结构、UDF 中的对象
执行内存(Execution Memory)可用内存 × fraction × (1 - storageFraction)Shuffle、Join、Sort、聚合等计算的临时数据
存储内存(Storage Memory)可用内存 × fraction × storageFractioncache/persist 的数据、广播变量

关键参数 ​

  • spark.memory.fraction:执行+存储内存占可用内存的比例,默认 0.6
  • spark.memory.storageFraction:存储内存占执行+存储总内存的比例,默认 0.5

💡 举个例子: Executor 内存 8GB(8192MB)

  • 预留:300MB
  • 可用:8192 - 300 = 7892MB
  • 执行+存储:7892 × 0.6 = 4735MB
    • 存储内存:4735 × 0.5 = 2367MB
    • 执行内存:4735 × 0.5 = 2367MB
  • 用户内存:7892 × 0.4 = 3157MB

5.2 统一内存管理:执行和存储可以互相借用 ​

Spark 1.6 之后用的是统一内存管理(Unified Memory Management)。

核心规则:

  1. 执行内存可以借用存储内存:如果存储内存没用满,执行内存可以借用
  2. 存储内存可以借用执行内存:如果执行内存没用满,存储内存可以借用
  3. 执行内存优先级更高:如果执行内存需要用了,可以把存储内存缓存的数据淘汰掉(如果是 MEMORY_AND_DISK 就溢写磁盘)
  4. 存储内存不能抢占执行内存:存储借了执行的内存,执行要用的时候,存储必须还
执行内存和存储内存的动态关系:

初始分配:
┌──────────────┬──────────────┐
│  执行 50%     │  存储 50%     │
└──────────────┴──────────────┘

存储没用满,执行借用:
┌───────────────────┬───────────┐
│  执行 70%(借了)  │  存储 30%  │
└───────────────────┴───────────┘

执行需要更多,抢占存储:
┌────────────────────────┬──────┐
│  执行 85%(抢占存储)    │ 15%  │  ← 存储的数据被淘汰/溢写
└────────────────────────┴──────┘

🎯 设计意图: 动态调整,提高内存利用率。 有的作业计算多(执行内存需求大),有的作业缓存多(存储内存需求大), 固定比例会浪费,动态借用更灵活。

5.3 常见内存问题与排查 ​

问题一:Executor OOM ​

可能原因:

  1. 数据倾斜,某个分区数据特别大
  2. 执行内存不够(Shuffle 数据量大)
  3. 用户内存不够(UDF 里创建了大对象)
  4. 缓存数据太多,占了太多内存

排查方法:

  1. 看 Spark UI → Executors 页,看哪个 Executor 挂了
  2. 看 Executor 日志,找 OOM 堆栈
  3. 看 Stage 详情,是不是数据倾斜(某个 Task 数据量特别大)

解决方法:

  1. 增加分区数,减小每个分区的数据量
  2. 处理数据倾斜
  3. 调大 Executor 内存
  4. 用 persist(MEMORY_AND_DISK),内存不够溢写磁盘
  5. 优化 UDF,不要创建大对象

问题二:GC 严重 ​

现象:Spark UI 上 Task 的 GC Time 占比很高(超过 10%)

原因:创建了太多短生命周期对象,JVM 频繁 GC

解决方法:

  1. 用 Kryo 序列化(比 Java 序列化快、省内存)
  2. 数据序列化存储(MEMORY_ONLY_SER)
  3. 优化代码,减少对象创建
  4. 调大 Executor 内存(给 GC 更多喘息空间)
  5. 用 G1 GC(spark.executor.extraJavaOptions=-XX:+UseG1GC)

六、调度机制:从 Job 到 Task 的完整链路 ​

6.1 调度的层级 ​

Spark 的调度分三层:

应用级别(Application)
    │  由集群管理器调度(YARN/K8s)
    ▼
Job 级别(Job)
    │  由 DAGScheduler 调度,拆成 Stage
    ▼
Stage 级别(Stage)
    │  由 TaskScheduler 调度,拆成 Task
    ▼
Task 级别(Task)
       由 Executor 的线程池执行

6.2 一个 Action 触发的完整调度流程 ​

1. 用户代码调用 Action(如 collect())
        │
        ▼
2. SparkContext.runJob() 提交 Job
        │
        ▼
3. DAGScheduler 接收 Job
   - 从最终 RDD 从后往前遍历依赖
   - 遇到宽依赖就切分 Stage
   - 生成 Stage 的依赖关系图(DAG of Stage)
   - 没有父 Stage 的 Stage 先提交
        │
        ▼
4. 每个 Stage 生成 TaskSet
   - Task 数量 = RDD 分区数
   - 每个 Task 处理一个分区
        │
        ▼
5. TaskScheduler 接收 TaskSet
   - 按调度策略(FIFO/FAIR)排队
   - 为每个 Task 计算数据本地性
        │
        ▼
6. SchedulerBackend 分配资源
   - 找空闲的 Executor
   - 考虑数据本地性(PROCESS_LOCAL > NODE_LOCAL > RACK_LOCAL > ANY)
        │
        ▼
7. Task 序列化后发到 Executor
        │
        ▼
8. Executor 反序列化 Task,在线程池中执行
        │
        ▼
9. Task 执行完,结果回传给 Driver
        │
        ▼
10. 所有 Task 完成 → Stage 完成 → 提交下一个 Stage
        │
        ▼
11. 所有 Stage 完成 → Job 完成 → 返回结果给用户

6.3 数据本地性(Data Locality) ​

Spark 会尽量把 Task 调度到数据所在的位置,减少网络传输。

本地性优先级从高到低:

级别含义说明
PROCESS_LOCAL同一个进程数据在当前 Executor 的内存里(缓存数据),最快
NODE_LOCAL同一个节点数据在当前节点的磁盘上(HDFS 本地块、Shuffle 数据)
RACK_LOCAL同一个机架数据在同机架的其他节点上,要走交换机
ANY任意位置数据在任意地方,要跨机架传输,最慢

本地性等待机制 ​

如果高本地性的节点暂时没有空闲资源,Spark 会等一会儿(spark.locality.wait,默认 3 秒),等不到就降级到下一级本地性。

💡 调优启示:

  • 如果数据本地性很差(很多 Task 是 ANY 级别),说明数据分布和计算资源不匹配
  • 可以考虑:让 Spark 节点和 HDFS 节点共部、增加副本数、调整等待时间

6.4 调度模式:FIFO vs FAIR ​

Spark 内部的 Task 调度有两种模式(spark.scheduler.mode):

模式特点适用场景
FIFO(默认)先提交的 Job 优先执行,前面的占满资源后面的等单用户、批处理
FAIR公平调度,多个 Job 轮流使用资源,支持权重多用户、多租户、交互式查询

FAIR 模式还支持配置调度池(Pool),给不同的池设置权重和最小资源,实现更精细的资源隔离。

🎯 生产环境建议:

  • 纯批处理、单用户 → FIFO(简单、 overhead 小)
  • 多用户共享集群、有交互式查询 → FAIR(避免大作业饿死小作业)

七、部署模式:Standalone / YARN / Mesos / K8s 深度对比 ​

7.1 客户端模式 vs 集群模式 ​

这是所有部署模式都有的一个概念:Driver 跑在哪?

模式Driver 位置特点适用场景
client(客户端模式)Driver 跑在提交任务的机器上日志直接在控制台看,方便调试;提交机器不能关开发调试、交互式
cluster(集群模式)Driver 跑在集群里(由集群管理器管理)Driver 由集群管理,有高可用;日志要去集群看生产环境

⚠️ 重要: client 模式下,如果你提交完任务就关了电脑/终端,Driver 就挂了,整个任务就失败了。 生产环境一定要用 cluster 模式。

7.2 四种集群管理器对比 ​

Standalone(Spark 自带) ​

架构:Master + Worker,Spark 自己管理资源

优点:

  • 简单,不依赖其他组件
  • 安装部署容易
  • 功能完整,支持所有 Spark 特性

缺点:

  • 只能跑 Spark,不能跑其他框架
  • 资源管理不如 YARN 成熟
  • 没有多租户、队列等高级功能

适用:学习测试、纯 Spark 的小集群

YARN(Hadoop 生态) ​

架构:ResourceManager + NodeManager,Spark 作为 YARN 的应用运行

优点:

  • 和 HDFS 天然集成(数据本地性好)
  • 资源调度成熟,支持队列、多租户
  • 生态完善,Hive、HBase、Flink 都能跑在 YARN 上
  • 生产环境最成熟稳定

缺点:

  • 依赖 Hadoop 生态,组件多
  • 弹性伸缩能力弱
  • 容器化支持不好

适用:传统大数据公司、已有 Hadoop 集群

Mesos ​

架构:Mesos Master + Agent,细粒度资源分配

优点:

  • 细粒度资源分配(可以按 CPU 核分配,不像 YARN 按 Container)
  • 支持多种框架(Spark、Marathon、Chronos...)
  • 资源利用率高

缺点:

  • 社区活跃度下降,用的人越来越少
  • 国内资料少
  • 学习成本高

适用:已有 Mesos 集群的环境(越来越少)

Kubernetes(云原生) ​

架构:Spark 作为 K8s 的 Pod 运行,Driver 和 Executor 都是 Pod

优点:

  • 云原生,和 K8s 生态无缝集成
  • 弹性伸缩好(可以根据负载动态增减 Executor)
  • 容器化,环境一致性好
  • 资源隔离更彻底
  • 支持 GPU 调度

缺点:

  • 相对较新,生产案例不如 YARN 多
  • 需要 K8s 运维能力
  • 和 HDFS 的数据本地性需要特殊配置

适用:云原生架构、互联网新公司、混合云

7.3 选型决策树 ​

你们有 Hadoop 集群吗?
├── 有 → 用 YARN(最稳妥,90% 的传统公司选这个)
└── 没有
    ├── 有 K8s 集群吗?
    │   ├── 有 → 用 Kubernetes(云原生方向)
    │   └── 没有
    │       ├── 只是学习/测试 → Standalone(最简单)
    │       └── 要上生产 → 建议搭 YARN 或 K8s
    └── 有 Mesos 集群吗?
        └── 有 → 用 Mesos(但建议考虑迁移到 K8s)

🎯 趋势判断: 短期(1~3年):YARN 仍然是生产环境的主流 长期(3~5年):Kubernetes 会越来越多,特别是新建的集群 Mesos 会逐渐退出历史舞台 Standalone 主要用于学习和测试


上篇小结 ​

到这里,你应该对 Spark 的核心原理有了深度理解:

  1. 定位:Spark 是第三代大数据计算引擎,解决了 MapReduce 慢、技术栈碎、开发效率低的问题,但不是银弹
  2. 架构:Driver + Executor 的主从架构,Driver 内部有 DAGScheduler、TaskScheduler 等组件
  3. RDD:本质是计算的抽象(不存数据,存计算描述),五大特性是理解 Spark 的钥匙
  4. 宽窄依赖:Spark 调度的灵魂,宽依赖触发 Shuffle,是性能瓶颈
  5. 内存管理:统一内存管理,执行和存储动态借用
  6. 调度:Job → Stage → Task 三层调度,数据本地性是关键优化点
  7. 部署:YARN 是当前主流,K8s 是未来方向

中篇预告:Spark SQL 深入(Catalyst 优化器、Tungsten 执行引擎)、Structured Streaming、MLlib、GraphX,以及 Spark vs Flink vs MapReduce 的深度对比。


文档版本:V1.0(上篇) 最后更新:2026 年 8 月 适用 Spark 版本:3.5.x

基于 Vite 强力驱动 | 纯静态轻量托管