Spark核心

8 minIntermediate2026/6/14

Spark核心概念:RDD、DataFrame、Spark SQL原理与实战。

1. Spark概述

Apache Spark 是一个统一的分布式计算引擎,基于内存计算提供比 MapReduce 快 10~100 倍的性能。

1.1 Spark核心特性

特性说明
内存计算中间结果缓存在内存,减少磁盘IO
DAG执行引擎优化执行计划,减少Shuffle
统一API批处理、流处理、ML、计算统一接口
多语言支持Scala、Java、Python、R
生态丰富Spark SQL、Streaming、MLlib、GraphX

1.2 Spark vs MapReduce

维度MapReduceSpark
计算模式磁盘中间结果内存中间结果
执行模型两阶段(Map/Reduce)DAG多阶段
迭代计算每次迭代写磁盘内存缓存复用
延迟分钟级秒级
编程模型Map/Reduce函数丰富的算子API

2. RDD弹性分布式数据集

RDD(Resilient Distributed Dataset)是 Spark 最基本的抽象,代表一个不可变、分区的、容错的分布式对象集合。

2.1 RDD核心属性

  1. 分区列表(Partitions):数据集的分片,每个分片对应一个计算任务
  2. 计算函数(Compute):每个分区的计算逻辑
  3. 依赖列表(Dependencies):RDD 之间的血缘关系
  4. 分区器(Partitioner):Key-Value 型 RDD 的分区策略
  5. 首选位置(Preferred Locations):分区计算的优先位置(数据本地性)

2.2 RDD创建方式

# 1. 从集合创建
rdd = sc.parallelize([1, 2, 3, 4, 5], numSlices=3)

# 2. 从外部存储创建
rdd = sc.textFile("hdfs://path/to/file")

# 3. 从其他RDD转换
rdd2 = rdd.map(lambda x: x * 2)

2.3 RDD操作

转换(Transformation)——懒执行,构建DAG:

算子说明Shuffle
map逐元素转换
filter过滤元素
flatMap展平转换
union合并RDD
groupByKey按Key分组
reduceByKey按Key聚合
join关联操作
distinct去重
repartition重新分区

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

算子说明
collect收集所有元素到Driver
count计数
reduce聚合
saveAsTextFile保存到文件
take(n)取前n个元素
foreach遍历每个元素

2.4 RDD血缘与容错

RDD 通过**血缘(Lineage)**实现容错,而非数据复制:

RDD_A ──map──→ RDD_B ──filter──→ RDD_C ──reduceByKey──→ RDD_D

                                              分区丢失时

                                              从RDD_C重新计算

窄依赖(Narrow Dependency) vs 宽依赖(Wide Dependency):

定义容错成本示例
窄依赖父分区最多被一个子分区使用低(重新计算单个分区)map、filter
宽依赖父分区被多个子分区使用高(需重新Shuffle)groupByKey、join

3. DataFrame与Dataset

3.1 DataFrame

DataFrame 是以命名列组织的分布式数据集,似于关系数据库中的表:

# 从JSON创建
df = spark.read.json("people.json")

# DataFrame操作
df.select("name", "age").filter(df.age > 20).groupBy("name").count()

# SQL查询
df.createOrReplaceTempView("people")
spark.sql("SELECT name, COUNT(*) FROM people WHERE age > 20 GROUP BY name")

3.2 Dataset

Dataset 是 DataFrame 的类型安全版本(Scala/Java):

case class Person(name: String, age: Long)
val ds: Dataset[Person] = spark.read.json("people.json").as[Person]
ds.filter(p => p.age > 20).map(p => p.name)

3.3 RDD vs DataFrame vs Dataset

维度RDDDataFrameDataset
型安全编译期运行期编译期
优化Catalyst + TungstenCatalyst + Tungsten
序列化Java/KryoTungsten二进制Tungsten二进制
GC开销
API函数式SQL/DSLSQL/DSL/函数式

4. Spark SQL与Catalyst优化器

4.1 Catalyst优化器

Catalyst 是 Spark SQL 的核心优化引擎,基于规则优化代价优化

Unresolved Logical Plan
        │  (解析)

Resolved Logical Plan
        │  (逻辑优化: 谓词下推、列裁剪、常量折叠)

Optimized Logical Plan
        │  (物理计划生成: 多策略选择)

Physical Plans
        │  (代价模型选择最优计划)

Selected Physical Plan
        │  (代码生成: Whole-Stage CodeGen)

RDDs Execution

4.2 常见优化规则

优化规则说明效果
谓词下推(Predicate Pushdown)将过滤条件尽早执行减少数据处理量
列裁剪(Column Pruning)只读取需要的列减少IO
常量折叠(Constant Folding)预计算常量表达式减少运行时计算
Join重排序小表驱动大表减少Shuffle数据量
Broadcast Join小表广播到所有节点避免Shuffle

4.3 Whole-Stage CodeGen

将多个物理算子融合为单个Java函数,消除虚函数调用:

// 无CodeGen: 虚函数调用
while (rows.hasNext()) {
    Row row = filter.next(rows.next());  // 虚函数调用
    if (row != null) {
        project.next(row);               // 虚函数调用
    }
}

// 有CodeGen: 融合为单函数
while (rows.hasNext()) {
    InternalRow row = rows.next();
    if (row.getInt(0) > 10) {            // 内联过滤
        result.setInt(0, row.getString(1)); // 内联投影
    }
}

5. Spark运行架构

5.1 核心组件

┌─────────────────────────────────────────┐
│              Driver                      │
│  SparkContext / SparkSession             │
│  DAGScheduler → TaskScheduler            │
└──────────────────┬──────────────────────┘

    ┌──────────────┼──────────────┐
    ▼              ▼              ▼
┌────────┐  ┌────────┐  ┌────────┐
│Executor│  │Executor│  │Executor│
│ Task1  │  │ Task3  │  │ Task5  │
│ Task2  │  │ Task4  │  │ Task6  │
│ Cache  │  │ Cache  │  │ Cache  │
└────────┘  └────────┘  └────────┘

5.2 作业执行层次

ApplicationJobStageTask\text{Application} \supset \text{Job} \supset \text{Stage} \supset \text{Task}

概念触发说明
Applicationspark-submit一个Spark应用程序
Jobaction算子一个行动操作触发一个Job
StageShuffle边界被Shuffle依赖划分的阶段
Task分区Stage中单个分区的计算任务

5.3 内存管理

Spark Executor 内存划分:

Executor Memory=Reserved+Unified(Storage+Execution)\text{Executor Memory} = \text{Reserved} + \text{Unified}(\text{Storage} + \text{Execution})

区域默认比例用途
Reserved Memory300MBSpark内部使用
Storage Memory50% of Unified缓存RDD/DataFrame
Execution Memory50% of UnifiedShuffle、排序、Join
User Memory未分配区域用户数据结构